From f7ae2d4dd5f08826606f50761e8d7dc6f5a7ce3b Mon Sep 17 00:00:00 2001 From: Peter Dodd Date: Thu, 13 Aug 2026 18:53:33 +0100 Subject: [PATCH] fix(redis2): remove orphaned directory index members on listing (#10735) * fix(redis2): remove orphaned directory index members on listing ListDirectoryEntries skipped index members whose value key was gone and left them in the ZSET, so the per-directory child index grew without bound under any TTL workload. Mirror the ZRem the logical-expiry branch already performs. Co-Authored-By: Claude Opus 5 (1M context) * fix(redis2): keep the index member when a concurrent insert recreates the value The orphan cleanup removed the member unconditionally, so an InsertEntry landing between FindEntry and the ZRem left a live value with no index member, invisible to listings until another InsertEntry on that path. UpdateEntry does not re-add it, so the loss persisted. Restore the member when the value is present again after the removal. The value key and the directory index key hash to different slots, so a Lua script or MULTI over both is not available to the cluster store. Co-Authored-By: Claude Opus 5 (1M context) --------- Co-authored-by: Claude Opus 5 (1M context) --- weed/filer/redis2/universal_redis_store.go | 17 +++ .../redis2/universal_redis_store_test.go | 143 ++++++++++++++++++ 2 files changed, 160 insertions(+) create mode 100644 weed/filer/redis2/universal_redis_store_test.go diff --git a/weed/filer/redis2/universal_redis_store.go b/weed/filer/redis2/universal_redis_store.go index 776952fc4..d4ce64d80 100644 --- a/weed/filer/redis2/universal_redis_store.go +++ b/weed/filer/redis2/universal_redis_store.go @@ -205,6 +205,7 @@ func (store *UniversalRedis2Store) ListDirectoryEntries(ctx context.Context, dir if err != nil { glog.V(0).InfofCtx(ctx, "list %s : %v", path, err) if err == filer_pb.ErrNotFound { + store.removeOrphanedDirectoryListMember(ctx, dirListKey, path, fileName) err = nil continue } @@ -233,6 +234,22 @@ func (store *UniversalRedis2Store) ListDirectoryEntries(ctx context.Context, dir return lastFileName, err } +func (store *UniversalRedis2Store) removeOrphanedDirectoryListMember(ctx context.Context, dirListKey string, path util.FullPath, fileName string) { + if err := store.Client.ZRem(ctx, dirListKey, fileName).Err(); err != nil { + return + } + + // InsertEntry writes the value before adding the member, so a value present + // again here may belong to an insert that found the member still in place + // and whose ZAddNX was therefore a no-op. + exists, err := store.Client.Exists(ctx, store.getKey(string(path))).Result() + if err == nil && exists == 0 { + return + } + + store.Client.ZAddNX(ctx, dirListKey, redis.Z{Score: 0, Member: fileName}) +} + func genDirectoryListKey(dir string) (dirList string) { return dir + DIR_LIST_MARKER } diff --git a/weed/filer/redis2/universal_redis_store_test.go b/weed/filer/redis2/universal_redis_store_test.go new file mode 100644 index 000000000..e5d17e5e6 --- /dev/null +++ b/weed/filer/redis2/universal_redis_store_test.go @@ -0,0 +1,143 @@ +package redis2 + +import ( + "context" + "fmt" + "os" + "testing" + "time" + + "github.com/redis/go-redis/v9" + + "github.com/seaweedfs/seaweedfs/weed/filer" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +func newTestStore(t *testing.T, keyPrefix string) (*UniversalRedis2Store, util.FullPath) { + t.Helper() + + if os.Getenv("RUN_REDIS_TESTS") != "1" { + t.Skip("redis2 tests are disabled. Start a redis-server and set RUN_REDIS_TESTS=1 to enable, REDIS_ADDR defaults to 127.0.0.1:6379.") + } + + addr := os.Getenv("REDIS_ADDR") + if addr == "" { + addr = "127.0.0.1:6379" + } + + ctx := context.Background() + client := redis.NewClient(&redis.Options{Addr: addr}) + if err := client.Ping(ctx).Err(); err != nil { + t.Fatalf("connect to redis at %s: %v", addr, err) + } + + store := &UniversalRedis2Store{Client: client, keyPrefix: keyPrefix} + store.loadSuperLargeDirectories(nil) + + dir := util.FullPath(fmt.Sprintf("/redis2_test_%d", time.Now().UnixNano())) + t.Cleanup(func() { + store.DeleteFolderChildren(ctx, dir) + store.DeleteEntry(ctx, dir) + client.Close() + }) + + return store, dir +} + +func insertTestEntry(t *testing.T, store *UniversalRedis2Store, path util.FullPath, ttlSec int32) { + t.Helper() + + now := time.Now() + if err := store.InsertEntry(context.Background(), &filer.Entry{ + FullPath: path, + Attr: filer.Attr{Crtime: now, Mtime: now, Mode: 0644, TtlSec: ttlSec}, + }); err != nil { + t.Fatalf("InsertEntry %s: %v", path, err) + } +} + +func listNames(t *testing.T, store *UniversalRedis2Store, dir util.FullPath) []string { + t.Helper() + + names := []string{} + if _, err := store.ListDirectoryEntries(context.Background(), dir, "", true, 100, func(entry *filer.Entry) (bool, error) { + _, name := entry.FullPath.DirAndName() + names = append(names, name) + return true, nil + }); err != nil { + t.Fatalf("ListDirectoryEntries %s: %v", dir, err) + } + return names +} + +func indexMembers(t *testing.T, store *UniversalRedis2Store, dir util.FullPath) []string { + t.Helper() + + members, err := store.Client.ZRangeByLex(context.Background(), store.getKey(genDirectoryListKey(string(dir))), &redis.ZRangeBy{Min: "-", Max: "+"}).Result() + if err != nil { + t.Fatalf("read directory index of %s: %v", dir, err) + } + return members +} + +func TestListDirectoryEntriesRemovesOrphanedIndexMembers(t *testing.T) { + for _, keyPrefix := range []string{"", "sw:"} { + t.Run("keyPrefix="+keyPrefix, func(t *testing.T) { + store, dir := newTestStore(t, keyPrefix) + + insertTestEntry(t, store, dir.Child("alive"), 0) + insertTestEntry(t, store, dir.Child("orphan"), 0) + + if err := store.Client.Del(context.Background(), store.getKey(string(dir.Child("orphan")))).Err(); err != nil { + t.Fatalf("drop value key: %v", err) + } + + if names := listNames(t, store, dir); len(names) != 1 || names[0] != "alive" { + t.Fatalf("listed %v, want [alive]", names) + } + + if members := indexMembers(t, store, dir); len(members) != 1 || members[0] != "alive" { + t.Fatalf("directory index holds %v, want [alive]", members) + } + }) + } +} + +func TestRemoveOrphanedDirectoryListMemberKeepsRecreatedEntry(t *testing.T) { + store, dir := newTestStore(t, "") + + path := dir.Child("recreated") + insertTestEntry(t, store, path, 0) + + store.removeOrphanedDirectoryListMember(context.Background(), store.getKey(genDirectoryListKey(string(dir))), path, "recreated") + + if members := indexMembers(t, store, dir); len(members) != 1 || members[0] != "recreated" { + t.Fatalf("directory index holds %v, want [recreated]", members) + } + + if names := listNames(t, store, dir); len(names) != 1 || names[0] != "recreated" { + t.Fatalf("listed %v, want [recreated]", names) + } +} + +func TestListDirectoryEntriesRemovesIndexMembersExpiredByRedis(t *testing.T) { + store, dir := newTestStore(t, "") + + insertTestEntry(t, store, dir.Child("ttl"), 1) + + time.Sleep(1500 * time.Millisecond) + + if exists, err := store.Client.Exists(context.Background(), store.getKey(string(dir.Child("ttl")))).Result(); err != nil { + t.Fatalf("check value key: %v", err) + } else if exists != 0 { + t.Fatal("redis did not expire the value key, the logical expiry path is not being bypassed") + } + + if names := listNames(t, store, dir); len(names) != 0 { + t.Fatalf("listed %v, want none", names) + } + + if members := indexMembers(t, store, dir); len(members) != 0 { + t.Fatalf("directory index holds %v, want none", members) + } +}