Skip to content
Merged
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
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -202,9 +202,9 @@ docker run --rm \
-listen :6479 \
-primary redis.internal:6379 \
-secondary elastickv.internal:6380 \
-elastickv-pool-size 4 \
-secondary-write-concurrency 2 \
-secondary-script-concurrency 1 \
-elastickv-pool-size 192 \
-secondary-write-concurrency 96 \
-secondary-script-concurrency 3 \
-mode dual-write
```

Expand Down
38 changes: 24 additions & 14 deletions adapter/dynamodb_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -268,22 +268,32 @@ func TestDynamoDB_LeaderHealthz(t *testing.T) {

for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
req, err := http.NewRequestWithContext(
context.Background(),
http.MethodGet,
"http://"+tc.addr+dynamoLeaderHealthPath,
nil,
)
require.NoError(t, err)
require.Eventually(t, func() bool {
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()

req, err := http.NewRequestWithContext(
ctx,
http.MethodGet,
"http://"+tc.addr+dynamoLeaderHealthPath,
nil,
)
if err != nil {
return false
}

resp, err := http.DefaultClient.Do(req)
require.NoError(t, err)
defer resp.Body.Close()
resp, err := http.DefaultClient.Do(req)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if err != nil {
return false
}
defer resp.Body.Close()

require.Equal(t, tc.status, resp.StatusCode)
body, err := io.ReadAll(resp.Body)
require.NoError(t, err)
require.Equal(t, tc.body, string(body))
body, err := io.ReadAll(resp.Body)
if err != nil {
return false
}
return resp.StatusCode == tc.status && string(body) == tc.body
}, leaderChurnRetryTimeout, leaderChurnRetryInterval)
})
}
}
Expand Down
9 changes: 8 additions & 1 deletion adapter/redis.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,8 +114,15 @@ const (
minKeyedArgs = 2
)

const cmdElasticKVZRemFast = "ELASTICKV.ZREMFAST"

const (
redisDispatchTimeout = 10 * time.Second
redisDispatchTimeout = 30 * time.Second
// redisLuaDispatchTimeout gives EVAL/EVALSHA enough room for migration
// scripts that expand into thousands of Redis calls. Regular commands stay
// on redisDispatchTimeout so non-script heavy paths cannot hold worker slots
// for the full script replay budget.
redisLuaDispatchTimeout = 5 * time.Minute
Comment thread
bootjp marked this conversation as resolved.
// defaultRedisBlockWaitFallback is the safety-net poll interval for
// blocking-command wait loops when no in-process write signal arrives.
// Signals cover normal XADD / ZADD / ZINCRBY wakeups immediately; this
Expand Down
5 changes: 5 additions & 0 deletions adapter/redis_command_specs.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,10 @@ var redisCommandSpecs = []redisCommandSpec{
{Constant: cmdDBSize, Name: "dbsize", Arity: 1, Flags: []string{redisCmdFlagReadonly}, FirstKey: 0, LastKey: 0, Step: 0},
{Constant: cmdDel, Name: "del", Arity: -2, Flags: []string{redisCmdFlagWrite}, FirstKey: 1, LastKey: -1, Step: 1},
{Constant: cmdDiscard, Name: "discard", Arity: 1, Flags: []string{redisCmdFlagAdmin}, FirstKey: 0, LastKey: 0, Step: 0},
// ELASTICKV.ZREMFAST is a secondary-replay-only ZREM variant for the
// Redis proxy. It skips the full wrong-type scan when the expected zset row
// is absent so stale BZPOP replays do not turn into wide collection probes.
{Constant: cmdElasticKVZRemFast, Name: "elastickv.zremfast", Arity: -3, Flags: []string{redisCmdFlagWrite}, FirstKey: 1, LastKey: 1, Step: 1},
Comment thread
bootjp marked this conversation as resolved.
{Constant: cmdEval, Name: "eval", Arity: -3, Flags: []string{redisCmdFlagWrite}, FirstKey: 0, LastKey: 0, Step: 0},
{Constant: cmdEvalSHA, Name: "evalsha", Arity: -3, Flags: []string{redisCmdFlagWrite}, FirstKey: 0, LastKey: 0, Step: 0},
{Constant: cmdExec, Name: "exec", Arity: 1, Flags: []string{redisCmdFlagAdmin}, FirstKey: 0, LastKey: 0, Step: 0},
Expand Down Expand Up @@ -254,6 +258,7 @@ func (r *RedisServer) buildRouteMap() map[string]func(redcon.Conn, redcon.Comman
cmdZRevRangeByScore: r.zrevrangebyscore,
cmdZScore: r.zscore,
}
handlers[cmdElasticKVZRemFast] = r.elasticKVZRemFast
m := make(map[string]func(redcon.Conn, redcon.Command), len(redisCommandSpecs))
for _, s := range redisCommandSpecs {
h, ok := handlers[s.Constant]
Expand Down
106 changes: 106 additions & 0 deletions adapter/redis_compat_helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -1112,6 +1112,112 @@ func (r *RedisServer) deleteLogicalKeyElems(ctx context.Context, key []byte, rea
return elems, existed, nil
}

func (r *RedisServer) deleteLogicalKeyElemsForType(ctx context.Context, key []byte, readTS uint64, typ redisValueType) ([]*kv.Elem[kv.OP], bool, error) {
switch typ {
case redisTypeNone:
return nil, false, nil
case redisTypeString:
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
return elems, true, err
case redisTypeList:
return r.deleteListLogicalKeyElems(ctx, key, readTS)
case redisTypeHash:
return r.deleteHashLogicalKeyElems(ctx, key, readTS)
case redisTypeSet:
return r.deleteSetLogicalKeyElems(ctx, key, readTS)
case redisTypeZSet:
return r.deleteZSetLogicalKeyElems(ctx, key, readTS)
case redisTypeStream:
return r.deleteStreamLogicalKeyElems(ctx, key, readTS)
}
return nil, false, errors.WithStack(errors.AssertionFailedf("unknown redis type %v", typ))
}

func (r *RedisServer) deleteStringLikeElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], error) {
var elems []*kv.Elem[kv.OP]
for _, internalKey := range [][]byte{
redisStrKey(key),
key, // legacy bare string key
redisHLLKey(key),
redisTTLKey(key),
} {
ok, err := r.store.ExistsAt(ctx, internalKey, readTS)
if err != nil {
return nil, errors.WithStack(err)
}
if ok {
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: internalKey})
}
}
return elems, nil
}

func (r *RedisServer) deleteListLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
listElems, err := r.deleteListElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, listElems...), true, nil
}

func (r *RedisServer) deleteHashLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisHashKey(key)})
hashElems, err := r.deleteWideColumnElems(ctx, readTS,
store.HashFieldScanPrefix(key), store.HashMetaKey(key), store.HashMetaDeltaScanPrefix(key))
if err != nil {
return nil, false, err
}
return append(elems, hashElems...), true, nil
}

func (r *RedisServer) deleteSetLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisSetKey(key)})
setElems, err := r.deleteWideColumnElems(ctx, readTS,
store.SetMemberScanPrefix(key), store.SetMetaKey(key), store.SetMetaDeltaScanPrefix(key))
if err != nil {
return nil, false, err
}
return append(elems, setElems...), true, nil
}

func (r *RedisServer) deleteZSetLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisZSetKey(key)})
zsetElems, err := r.deleteZSetWideColumnElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, zsetElems...), true, nil
}

func (r *RedisServer) deleteStreamLogicalKeyElems(ctx context.Context, key []byte, readTS uint64) ([]*kv.Elem[kv.OP], bool, error) {
elems, err := r.deleteStringLikeElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
elems = append(elems, &kv.Elem[kv.OP]{Op: kv.Del, Key: redisStreamKey(key)})
streamElems, err := r.deleteStreamWideColumnElems(ctx, key, readTS)
if err != nil {
return nil, false, err
}
return append(elems, streamElems...), true, nil
}

// deleteStreamWideColumnElems returns delete operations for all stream
// wide-column keys: the meta key (if it exists) and every entry under the
// entry scan prefix. Total results are capped at maxWideColumnItems to
Expand Down
2 changes: 1 addition & 1 deletion adapter/redis_lua.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ func (r *RedisServer) runLuaScript(conn redcon.Conn, script string, evalArgs [][
return
}

ctx, cancel := context.WithTimeout(r.handlerContext(), redisDispatchTimeout)
ctx, cancel := context.WithTimeout(r.handlerContext(), redisLuaDispatchTimeout)
defer cancel()

start := time.Now()
Expand Down
19 changes: 14 additions & 5 deletions adapter/redis_lua_context.go
Original file line number Diff line number Diff line change
Expand Up @@ -3534,7 +3534,7 @@ func (c *luaScriptContext) commit() error {
}
sort.Strings(keys)

ctx, cancel := context.WithTimeout(c.scriptCtx(), redisDispatchTimeout)
ctx, cancel := context.WithTimeout(c.scriptCtx(), redisLuaDispatchTimeout)
defer cancel()

// Pre-allocate a commitTS so Delta key bytes can embed it before dispatch.
Expand Down Expand Up @@ -3658,18 +3658,27 @@ func (c *luaScriptContext) commitPlanForKey(ctx context.Context, key string, com
return luaKeyPlan{}, err
}

startType, err := c.server.keyTypeAt(ctx, []byte(key), c.startTS)
keyBytes := []byte(key)
rawStartType, err := c.server.rawKeyTypeAt(ctx, keyBytes, c.startTS)
if err != nil {
return luaKeyPlan{}, err
}
startType, err := c.server.applyTTLFilter(ctx, keyBytes, c.startTS, rawStartType)
if err != nil {
return luaKeyPlan{}, err
}
var deleteElems []*kv.Elem[kv.OP]
readKeys := luaWideFenceReadKeysForPlan([]byte(key), finalType, startType, valuePlan.preserveExisting)
readKeys := luaWideFenceReadKeysForPlan(keyBytes, finalType, startType, valuePlan.preserveExisting)
if !valuePlan.preserveExisting {
deleteElems, _, err = c.server.deleteLogicalKeyElems(ctx, []byte(key), c.startTS)
if c.everDeleted[key] && rawStartType != redisTypeNone {
deleteElems, _, err = c.server.deleteLogicalKeyElems(ctx, keyBytes, c.startTS)
} else {
deleteElems, _, err = c.server.deleteLogicalKeyElemsForType(ctx, keyBytes, c.startTS, rawStartType)
}
if err != nil {
return luaKeyPlan{}, err
}
deleteElems = append(deleteElems, redisTxnWideCollectionFenceElems([]byte(key))...)
deleteElems = append(deleteElems, redisTxnWideCollectionFenceElems(keyBytes)...)
}

dataElems := make([]*kv.Elem[kv.OP], 0, len(deleteElems)+len(valuePlan.elems))
Expand Down
7 changes: 4 additions & 3 deletions adapter/redis_peer_limiter.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ import (

const (
redisPerPeerLimitEnv = "ELASTICKV_REDIS_PER_PEER_CONNECTIONS"
defaultRedisProxyPoolPeerCap = 64
defaultRedisDedicatedPeerHeadroom = 64
defaultRedisPerPeerConnectionCap = defaultRedisProxyPoolPeerCap + defaultRedisDedicatedPeerHeadroom
defaultRedisProxyPoolPeerCap = 192
defaultRedisProxyReplicasPerPeer = 2
defaultRedisDedicatedPeerHeadroom = 128
defaultRedisPerPeerConnectionCap = defaultRedisProxyPoolPeerCap*defaultRedisProxyReplicasPerPeer + defaultRedisDedicatedPeerHeadroom
redisPeerLimitError = "ERR max connections per client exceeded"
unknownRedisPeer = "unknown"
)
Expand Down
6 changes: 5 additions & 1 deletion adapter/redis_peer_limiter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,11 @@ func TestRedisPeerLimiterDefaultMatchesProxyPool(t *testing.T) {
t.Setenv(redisPerPeerLimitEnv, "")
limiter := newDefaultRedisPeerLimiter()
require.NotNil(t, limiter)
require.Equal(t, proxy.DefaultElasticKVBackendOptions().PoolSize+defaultRedisDedicatedPeerHeadroom, limiter.limit)
require.Equal(t, proxy.DefaultElasticKVBackendOptions().PoolSize, defaultRedisProxyPoolPeerCap)
require.Equal(t,
proxy.DefaultElasticKVBackendOptions().PoolSize*defaultRedisProxyReplicasPerPeer+defaultRedisDedicatedPeerHeadroom,
limiter.limit,
)
}

func TestRedisPeerLimiterRejectsAndReleases(t *testing.T) {
Expand Down
53 changes: 53 additions & 0 deletions adapter/redis_retry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -749,3 +749,56 @@ func TestZRemDeletesWideColumnRows(t *testing.T) {
require.True(t, exists)
require.Equal(t, []redisZSetEntry{{Member: "b", Score: 2.0}}, zset.Entries)
}

func TestElasticKVZRemFastSkipsWrongTypeMiss(t *testing.T) {
t.Parallel()

ctx := context.Background()
st := store.NewMVCCStore()
key := []byte("internal:zremfast:string")
require.NoError(t, st.PutAt(ctx, redisStrKey(key), encodeRedisStr([]byte("value"), nil), 1, 0))

coord := newRetryOnceCoordinator(st)
coord.clock.Observe(1)
srv := NewRedisServer(nil, "", st, coord, nil, nil)

normalConn := &recordingConn{}
srv.zrem(normalConn, redcon.Command{Args: [][]byte{[]byte("ZREM"), key, []byte("member")}})
require.Contains(t, normalConn.err, wrongTypeMessage)

fastConn := &recordingConn{}
srv.elasticKVZRemFast(fastConn, redcon.Command{Args: [][]byte{[]byte(cmdElasticKVZRemFast), key, []byte("member")}})
require.Empty(t, fastConn.err)
require.Equal(t, int64(0), fastConn.int)
}

func TestElasticKVZRemFastRemovesLegacyBlobZSet(t *testing.T) {
t.Parallel()

ctx := context.Background()
st := store.NewMVCCStore()
key := []byte("internal:zremfast:legacy")
payload, err := marshalZSetValue(redisZSetValue{
Entries: []redisZSetEntry{
{Member: "a", Score: 1},
{Member: "b", Score: 2},
},
})
require.NoError(t, err)
require.NoError(t, st.PutAt(ctx, redisZSetKey(key), payload, 1, 0))

coord := newRetryOnceCoordinator(st)
coord.clock.Observe(1)
srv := NewRedisServer(nil, "", st, coord, nil, nil)

conn := &recordingConn{}
srv.elasticKVZRemFast(conn, redcon.Command{Args: [][]byte{[]byte(cmdElasticKVZRemFast), key, []byte("a")}})

require.Empty(t, conn.err)
require.Equal(t, int64(1), conn.int)

zset, exists, err := srv.loadZSetAt(ctx, key, snapshotTS(coord.clock, st))
require.NoError(t, err)
require.True(t, exists)
require.Equal(t, []redisZSetEntry{{Member: "b", Score: 2}}, zset.Entries)
}
Loading
Loading