diff --git a/state/redis/redis.go b/state/redis/redis.go index 3acf712fcc..86a26805e2 100644 --- a/state/redis/redis.go +++ b/state/redis/redis.go @@ -459,7 +459,9 @@ func (r *StateStore) registerSchemas(ctx context.Context) error { for name, elem := range r.querySchemas { r.logger.Infof("create query index %s", name) if err := r.client.DoWrite(ctx, elem.schema...); err != nil { - if err.Error() != "Index already exists" { + // Redis < 8.8 replies "Index already exists"; Redis >= 8.8 prefixes + // the reply with an error code: "SEARCH_INDEX_EXISTS Index already exists". + if !strings.Contains(err.Error(), "Index already exists") { return err } r.logger.Infof("drop stale query index %s", name) diff --git a/state/redis/redis_test.go b/state/redis/redis_test.go index 6708a78835..5179a5b43d 100644 --- a/state/redis/redis_test.go +++ b/state/redis/redis_test.go @@ -14,6 +14,8 @@ limitations under the License. package redis import ( + "context" + "errors" "strconv" "testing" "time" @@ -585,3 +587,60 @@ func Test_KeyList(t *testing.T) { _, ok := s.(state.KeysLiker) require.True(t, ok) } + +type fakeQueryIndexClient struct { + rediscomponent.RedisClient + createErr error + calls []string +} + +func (c *fakeQueryIndexClient) DoWrite(ctx context.Context, args ...interface{}) error { + cmd, _ := args[0].(string) + c.calls = append(c.calls, cmd) + if cmd == "FT.CREATE" { + err := c.createErr + // The index is dropped before the retry, so the retry succeeds. + c.createErr = nil + return err + } + return nil +} + +func TestRegisterSchemas(t *testing.T) { + newStore := func(t *testing.T, client rediscomponent.RedisClient) *StateStore { + store := newStateStore(logger.NewLogger("test")) + schemas, err := parseQuerySchemas(`[{"name":"userIdx","indexes":[{"key":"user.email","type":"TEXT"}]}]`) + require.NoError(t, err) + store.querySchemas = schemas + store.client = client + return store + } + + t.Run("index does not exist", func(t *testing.T) { + client := &fakeQueryIndexClient{} + store := newStore(t, client) + require.NoError(t, store.registerSchemas(t.Context())) + assert.Equal(t, []string{"FT.CREATE"}, client.calls) + }) + + t.Run("index exists on Redis < 8.8", func(t *testing.T) { + client := &fakeQueryIndexClient{createErr: errors.New("Index already exists")} + store := newStore(t, client) + require.NoError(t, store.registerSchemas(t.Context())) + assert.Equal(t, []string{"FT.CREATE", "FT.DROPINDEX", "FT.CREATE"}, client.calls) + }) + + t.Run("index exists on Redis >= 8.8", func(t *testing.T) { + client := &fakeQueryIndexClient{createErr: errors.New("SEARCH_INDEX_EXISTS Index already exists")} + store := newStore(t, client) + require.NoError(t, store.registerSchemas(t.Context())) + assert.Equal(t, []string{"FT.CREATE", "FT.DROPINDEX", "FT.CREATE"}, client.calls) + }) + + t.Run("unrelated error is returned", func(t *testing.T) { + client := &fakeQueryIndexClient{createErr: errors.New("LOADING Redis is loading the dataset in memory")} + store := newStore(t, client) + require.Error(t, store.registerSchemas(t.Context())) + assert.Equal(t, []string{"FT.CREATE"}, client.calls) + }) +}