mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-09 00:56:37 +08:00
refactor(oauth): replace legacy oauth cache with standard ram cache and add pubsub synchronization
- Replaced custom map-based cache in apps/oauth/cache.go with standard pkg/cache/ram framework. - Implemented Redis Pub/Sub invalidation channels for distributed token and user cache synchronization. - Created apps/oauth/cache_test.go to verify local cache operations and pub/sub broadcasts. refactor(cache): generic RAM cache with CoW and unified preheating Replaced L2 Redis cache and old cache package with process-local generic pkg/cache/ram. Implemented Copy-on-Write for reads, fine-grained locks per type for writes, and unified preheating in bootstrap. Changed cache invalidation to lazy-loading to resolve SQLite deadlocks during transactions.
This commit is contained in:
@@ -250,12 +250,14 @@ func TestLoadUsesEnvConfigWhenFileIsMissing(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestLoadDetectsOutboundIPWhenNodeIPMissing(t *testing.T) {
|
||||
nodeip.ResetCacheForTest()
|
||||
previousLookup := nodeip.LookupOutboundIP
|
||||
nodeip.LookupOutboundIP = func(ctx context.Context, strategies ...geoip.OutboundIPStrategy) (net.IP, error) {
|
||||
return net.ParseIP("8.8.8.8"), nil
|
||||
}
|
||||
defer func() {
|
||||
nodeip.LookupOutboundIP = previousLookup
|
||||
nodeip.ResetCacheForTest()
|
||||
}()
|
||||
|
||||
dir := t.TempDir()
|
||||
@@ -283,6 +285,7 @@ func TestLoadDetectsOutboundIPWhenNodeIPMissing(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestLoadFallsBackToLocalIPWhenOutboundLookupFails(t *testing.T) {
|
||||
nodeip.ResetCacheForTest()
|
||||
previousOutboundLookup := nodeip.LookupOutboundIP
|
||||
previousLocalLookup := nodeip.LookupLocalIP
|
||||
nodeip.LookupOutboundIP = func(ctx context.Context, strategies ...geoip.OutboundIPStrategy) (net.IP, error) {
|
||||
@@ -294,6 +297,7 @@ func TestLoadFallsBackToLocalIPWhenOutboundLookupFails(t *testing.T) {
|
||||
defer func() {
|
||||
nodeip.LookupOutboundIP = previousOutboundLookup
|
||||
nodeip.LookupLocalIP = previousLocalLookup
|
||||
nodeip.ResetCacheForTest()
|
||||
}()
|
||||
|
||||
dir := t.TempDir()
|
||||
|
||||
@@ -7,7 +7,6 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
@@ -46,13 +45,8 @@ func TestListVisibleSystemConfigsUsesRedisCache(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
exists, err := db.Redis.Exists(ctx, db.PrefixedKey(repository.SystemConfigVisibleListRedisKey)).Result()
|
||||
if err != nil {
|
||||
t.Fatalf("Redis.Exists() error = %v", err)
|
||||
}
|
||||
if exists == 0 {
|
||||
t.Fatal("expected visible config list cache key to exist")
|
||||
}
|
||||
// Since system configs are now purely cached in process-local RAM (L1) and not written to Redis (L2),
|
||||
// we do not verify the existence of the legacy Redis key here.
|
||||
|
||||
if err := repository.InvalidateVisibleSystemConfigsCache(ctx); err != nil {
|
||||
t.Fatalf("InvalidateVisibleSystemConfigsCache() error = %v", err)
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
|
||||
@@ -25,6 +26,7 @@ func TestSystemConfigRAMCacheServesUntilInvalidated(t *testing.T) {
|
||||
if err := repository.InvalidateAllSystemConfigCaches(ctx); err != nil {
|
||||
t.Fatalf("InvalidateAllSystemConfigCaches() error = %v", err)
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
warm, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeySiteName)
|
||||
if err != nil {
|
||||
@@ -54,6 +56,7 @@ func TestSystemConfigRAMCacheServesUntilInvalidated(t *testing.T) {
|
||||
if err := repository.InvalidateSystemConfigCache(ctx, model.ConfigKeySiteName); err != nil {
|
||||
t.Fatalf("InvalidateSystemConfigCache(site_name) error = %v", err)
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
refreshed, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeySiteName)
|
||||
if err != nil {
|
||||
@@ -63,13 +66,8 @@ func TestSystemConfigRAMCacheServesUntilInvalidated(t *testing.T) {
|
||||
t.Fatalf("GetSystemConfigByKey(site_name).Value = %q, want %q", refreshed.Value, "ram_probe_value")
|
||||
}
|
||||
|
||||
exists, err := db.Redis.HExists(ctx, db.PrefixedKey(repository.SystemConfigRedisHashKey), model.ConfigKeySiteName).Result()
|
||||
if err != nil {
|
||||
t.Fatalf("HExists(site_name) error = %v", err)
|
||||
}
|
||||
if !exists {
|
||||
t.Fatal("expected redis hash field to be repopulated after refresh")
|
||||
}
|
||||
// Since system configs are now purely cached in process-local RAM (L1) and not written to Redis (L2),
|
||||
// we do not verify if the Redis hash field is repopulated.
|
||||
}
|
||||
|
||||
func TestInvalidateSystemConfigCacheClearsRedisField(t *testing.T) {
|
||||
@@ -86,6 +84,7 @@ func TestInvalidateSystemConfigCacheClearsRedisField(t *testing.T) {
|
||||
if err := repository.InvalidateSystemConfigCache(ctx, model.ConfigKeySiteName); err != nil {
|
||||
t.Fatalf("InvalidateSystemConfigCache(site_name) error = %v", err)
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
_, err = db.Redis.HGet(ctx, db.PrefixedKey(repository.SystemConfigRedisHashKey), model.ConfigKeySiteName).Result()
|
||||
if !errors.Is(err, redis.Nil) {
|
||||
|
||||
@@ -104,3 +104,11 @@ func DetectLocal() string {
|
||||
}
|
||||
return bestIP
|
||||
}
|
||||
|
||||
// ResetCacheForTest clears the cached IP.
|
||||
func ResetCacheForTest() {
|
||||
cacheMu.Lock()
|
||||
cachedIP = ""
|
||||
lastDetected = time.Time{}
|
||||
cacheMu.Unlock()
|
||||
}
|
||||
|
||||
@@ -115,3 +115,17 @@ func clearSystemConfigCache() {
|
||||
log.Printf("[%s] clear system config cache failed: %v\n", dbType(), err)
|
||||
}
|
||||
}
|
||||
|
||||
func tableExistsSQL(dialect string) string {
|
||||
if dialect == dialectPostgres {
|
||||
return "SELECT count(*) FROM information_schema.tables WHERE table_schema='public' AND table_name=$1"
|
||||
}
|
||||
return "SELECT count(*) FROM sqlite_master WHERE type='table' AND name=?"
|
||||
}
|
||||
|
||||
func tablesWithPrefixSQL(dialect string) string {
|
||||
if dialect == dialectPostgres {
|
||||
return "SELECT table_name FROM information_schema.tables WHERE table_schema='public' AND table_name LIKE $1"
|
||||
}
|
||||
return "SELECT name FROM sqlite_master WHERE type='table' AND name LIKE ?"
|
||||
}
|
||||
|
||||
@@ -5,10 +5,13 @@ package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/pkg/cache/ram"
|
||||
"github.com/alicebob/miniredis/v2"
|
||||
"github.com/glebarez/sqlite"
|
||||
"github.com/redis/go-redis/v9"
|
||||
@@ -76,7 +79,7 @@ func TestListSystemConfigsByKeys_EmptyKeys(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestListSystemConfigsByKeys_LoadsFromRedisBeforeDB(t *testing.T) {
|
||||
func TestListSystemConfigsByKeys_LoadsFromRAMBeforeDB(t *testing.T) {
|
||||
dbConn, cleanup := setupSystemConfigTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
@@ -85,6 +88,7 @@ func TestListSystemConfigsByKeys_LoadsFromRedisBeforeDB(t *testing.T) {
|
||||
if err := InvalidateAllSystemConfigCaches(ctx); err != nil {
|
||||
t.Fatalf("InvalidateAllSystemConfigCaches() error = %v", err)
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
warm, err := GetSystemConfigByKey(ctx, model.ConfigKeySiteName)
|
||||
if err != nil {
|
||||
@@ -100,8 +104,6 @@ func TestListSystemConfigsByKeys_LoadsFromRedisBeforeDB(t *testing.T) {
|
||||
t.Fatalf("Update(site_name) error = %v", err)
|
||||
}
|
||||
|
||||
ResetSystemConfigRAMCacheForTest()
|
||||
|
||||
configs, err := ListSystemConfigsByKeys(ctx, []string{model.ConfigKeySiteName})
|
||||
if err != nil {
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name) error = %v", err)
|
||||
@@ -112,11 +114,28 @@ func TestListSystemConfigsByKeys_LoadsFromRedisBeforeDB(t *testing.T) {
|
||||
t.Fatal("ListSystemConfigsByKeys(site_name) missing site_name entry")
|
||||
}
|
||||
if sc.Value != "Wavelet" {
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name).Value = %q, want redis value %q", sc.Value, "Wavelet")
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name).Value = %q, want cached RAM value %q", sc.Value, "Wavelet")
|
||||
}
|
||||
|
||||
if err := InvalidateAllSystemConfigCaches(ctx); err != nil {
|
||||
t.Fatalf("InvalidateAllSystemConfigCaches() error = %v", err)
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
configs, err = ListSystemConfigsByKeys(ctx, []string{model.ConfigKeySiteName})
|
||||
if err != nil {
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name) error = %v", err)
|
||||
}
|
||||
sc, ok = configs[model.ConfigKeySiteName]
|
||||
if !ok {
|
||||
t.Fatal("ListSystemConfigsByKeys(site_name) missing site_name entry after invalidate")
|
||||
}
|
||||
if sc.Value != "db_only_value" {
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name).Value = %q, want db value %q", sc.Value, "db_only_value")
|
||||
}
|
||||
}
|
||||
|
||||
func TestListSystemConfigsByKeys_PopulatesRAMFromRedis(t *testing.T) {
|
||||
func TestListSystemConfigsByKeys_PopulatesRAMOnMiss(t *testing.T) {
|
||||
_, cleanup := setupSystemConfigTest(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
@@ -125,22 +144,21 @@ func TestListSystemConfigsByKeys_PopulatesRAMFromRedis(t *testing.T) {
|
||||
if err := InvalidateAllSystemConfigCaches(ctx); err != nil {
|
||||
t.Fatalf("InvalidateAllSystemConfigCaches() error = %v", err)
|
||||
}
|
||||
|
||||
if _, err := GetSystemConfigByKey(ctx, model.ConfigKeySiteName); err != nil {
|
||||
t.Fatalf("GetSystemConfigByKey(site_name) warm error = %v", err)
|
||||
}
|
||||
|
||||
ResetSystemConfigRAMCacheForTest()
|
||||
time.Sleep(50 * time.Millisecond) // Wait for async Redis broadcast to be processed
|
||||
|
||||
if _, err := ListSystemConfigsByKeys(ctx, []string{model.ConfigKeySiteName}); err != nil {
|
||||
t.Fatalf("ListSystemConfigsByKeys(site_name) error = %v", err)
|
||||
}
|
||||
|
||||
cached, ok := systemConfigRAMCache.GetIfPresent(model.ConfigKeySiteName)
|
||||
cachedItem, ok := ram.Get(ConfigCacheType, model.ConfigKeySiteName)
|
||||
if !ok {
|
||||
t.Fatal("expected RAM cache to be populated after redis hit")
|
||||
t.Fatal("expected RAM cache to be populated after config query")
|
||||
}
|
||||
if cached.Value != "Wavelet" {
|
||||
t.Fatalf("RAM cache value = %q, want %q", cached.Value, "Wavelet")
|
||||
var cachedConfig model.SystemConfig
|
||||
if err := json.Unmarshal([]byte(cachedItem.Value), &cachedConfig); err != nil {
|
||||
t.Fatalf("unmarshal cached value error = %v", err)
|
||||
}
|
||||
}
|
||||
if cachedConfig.Value != "Wavelet" {
|
||||
t.Fatalf("RAM cache value = %q, want %q", cachedConfig.Value, "Wavelet")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user