Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion state/redis/redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
59 changes: 59 additions & 0 deletions state/redis/redis_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ limitations under the License.
package redis

import (
"context"
"errors"
"strconv"
"testing"
"time"
Expand Down Expand Up @@ -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)
})
}
Loading