Skip to content
Merged
Show file tree
Hide file tree
Changes from 10 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
53 changes: 53 additions & 0 deletions adapter/redis_error_prefix_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,33 @@ func TestHandleProxyTxnError(t *testing.T) {

}

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

for _, tc := range []struct {
name string
err error
}{
{name: "typed", err: errors.WithStack(errRedisExecRouteChangedAfterAmbiguousAttempt)},
{name: "decoded RESP", err: errors.New(errRedisExecRouteChangedAfterAmbiguousAttempt.Error())},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
c := &captureConn{}
handled := handleProxyTxnError(c, tc.err)
if !handled {
t.Fatal("handleProxyTxnError returned false")
}
if c.lastErr != errRedisExecRouteChangedAfterAmbiguousAttempt.Error() {
t.Fatalf("last error = %q", c.lastErr)
}
if c.wroteArray {
t.Fatalf("unexpected array reply %d", c.lastArray)
}
})
}
}

func TestHandleProxyTxnHeavyCommandBusyError(t *testing.T) {
t.Parallel()
c := &captureConn{}
Expand Down Expand Up @@ -222,6 +249,32 @@ func TestHandleProxyTxnCommandError(t *testing.T) {
}
})

t.Run("ambiguous route change is promoted to top-level EXEC error", func(t *testing.T) {
t.Parallel()

for _, tc := range []struct {
name string
err error
}{
{name: "typed", err: errors.WithStack(errRedisExecRouteChangedAfterAmbiguousAttempt)},
{name: "decoded RESP", err: errors.New(errRedisExecRouteChangedAfterAmbiguousAttempt.Error())},
} {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
cmd := redis.NewCmd(context.Background(), "SET", "k", "v")
cmd.SetErr(tc.err)
c := &captureConn{}
handled := handleProxyTxnCommandError(c, []*redis.Cmd{cmd})
if !handled {
t.Fatal("handleProxyTxnCommandError returned false")
}
if c.lastErr != errRedisExecRouteChangedAfterAmbiguousAttempt.Error() {
t.Fatalf("last error = %q", c.lastErr)
}
})
}
})

}

func TestHandleProxyTxnCommandHeavyCommandBusyError(t *testing.T) {
Expand Down
64 changes: 51 additions & 13 deletions adapter/redis_lists.go
Original file line number Diff line number Diff line change
Expand Up @@ -265,7 +265,7 @@ func (r *RedisServer) dispatchListPushReuse(ctx context.Context, key []byte, pen
// (non-retryable errors escape to the client; pending is then
// discarded with the goroutine, so the update is wasted and the
// stale value would be misleading if some future caller reads it).
if isRetryableRedisTxnErr(dispErr) {
if isReusableRedisTxnErr(dispErr) {
pending.commitTS = commitTS
}
return 0, false, errors.WithStack(dispErr)
Expand Down Expand Up @@ -440,7 +440,7 @@ func (r *RedisServer) listPushCoreWithDedup(ctx context.Context, key []byte, val
// retryRedisWrite's retry predicate; ambiguous errors that escape
// to the client are a separate problem space (cross-request
// idempotency cache) and out of scope for this design.
if isRetryableRedisTxnErr(dispErr) {
if isReusableRedisTxnErr(dispErr) {
pending = &reusableListPush{
ops: ops,
startTS: startTS,
Expand Down Expand Up @@ -709,11 +709,40 @@ func (r *RedisServer) fetchListRange(ctx context.Context, key []byte, meta store
}

func (r *RedisServer) rangeList(ctx context.Context, key []byte, startRaw, endRaw []byte) ([]string, error) {
if !r.coordinator.IsLeaderForKey(key) {
return r.proxyLRange(key, startRaw, endRaw)
var out []string
err := r.retryRedisWrite(ctx, func() error {
routeVersion := r.redisReadFenceRouteVersion()
readTS, readPin, proxied, ok, err := r.fenceRangeListReadGroups(ctx, key, startRaw, endRaw)
if err != nil {
return err
} else if ok {
if err := r.ensureRedisReadFenceRouteStable(routeVersion); err != nil {
return err
}
out = proxied
return nil
}
defer readPin.Release()
if err := r.ensureRedisReadFenceRouteStable(routeVersion); err != nil {
return err
}
next, err := r.rangeListAt(ctx, key, startRaw, endRaw, readTS)
if err != nil {
return err
}
if err := r.ensureRedisReadFenceRouteStable(routeVersion); err != nil {
return err
}
out = next
return nil
})
if err != nil {
return nil, err
}
return out, nil
}

readTS := r.readTS()
func (r *RedisServer) rangeListAt(ctx context.Context, key []byte, startRaw, endRaw []byte, readTS uint64) ([]string, error) {
typ, err := r.keyTypeAt(ctx, key, readTS)
if err != nil {
return nil, err
Expand All @@ -725,14 +754,6 @@ func (r *RedisServer) rangeList(ctx context.Context, key []byte, startRaw, endRa
return nil, wrongTypeError()
}

// PR #749 follow-up: pass the per-call dispatch ctx so a stalled
// VerifyLeaderForKey honours the caller's deadline rather than the
// long-lived handlerContext + verifyLeaderEngineCtx fallback. Same
// shape as keys() / FLUSHDB.
if err := r.coordinator.VerifyLeaderForKey(ctx, key); err != nil {
return nil, errors.WithStack(err)
}

meta, exists, err := r.resolveListMeta(ctx, key, readTS)
if err != nil {
return nil, err
Expand All @@ -749,6 +770,23 @@ func (r *RedisServer) rangeList(ctx context.Context, key []byte, startRaw, endRa
return r.fetchListRange(ctx, key, meta, int64(s), int64(e), readTS)
}

func (r *RedisServer) fenceRangeListReadGroups(ctx context.Context, key []byte, startRaw, endRaw []byte) (uint64, *kv.ActiveTimestampToken, []string, bool, error) {
groupKeys := r.redisReadFenceGroupKeys(r.redisTxnReadFenceKeysForRanges(key, redisListReadFenceRanges(key)))
proxyKey, ok, err := r.readFenceProxyKey(groupKeys)
if err != nil {
return 0, nil, nil, false, err
}
if ok {
proxied, err := r.proxyLRange(key, proxyKey, startRaw, endRaw)
return 0, nil, proxied, true, err
}
readTS, readPin, err := r.redisReadFencedTimestamp(ctx, groupKeys, r.readTS)
if err != nil {
return 0, nil, nil, false, err
}
return readTS, readPin, nil, false, nil
}

type listPushFunc func(ctx context.Context, key []byte, values [][]byte) (int64, error)
type listProxyFunc func(key []byte, values [][]byte) (int64, error)

Expand Down
59 changes: 50 additions & 9 deletions adapter/redis_proxy_leader.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,18 @@ func (r *RedisServer) proxyTransactionToLeader(conn redcon.Conn, queue []redcon.
if !ok {
return
}
r.proxyTransactionToLeaderAddr(conn, queue, leaderAddr)
}

func (r *RedisServer) proxyTransactionToLeaderForKey(conn redcon.Conn, routingKey []byte, queue []redcon.Command) {
leaderAddr, ok := r.resolveLeaderRedisAddrForKey(conn, routingKey)
if !ok {
return
}
r.proxyTransactionToLeaderAddr(conn, queue, leaderAddr)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve terminal EXEC errors across the proxy

When this new per-shard proxy path targets a node whose EXEC returns the newly introduced errRedisExecRouteChangedAfterAmbiguousAttempt, go-redis reports the top-level non-array EXEC error as the pipeline error and copies it onto the queued command handles. Neither proxy error handler recognizes this error, so proxyTransactionToLeaderAddr writes an EXEC result array containing per-command errors instead of preserving the top-level ambiguous-transaction failure; clients can therefore treat the transaction as having executed with ordinary command failures. Promote this EXEC-level error in the proxy handlers before writing the result array.

Useful? React with 👍 / 👎.

}

func (r *RedisServer) proxyTransactionToLeaderAddr(conn redcon.Conn, queue []redcon.Command, leaderAddr string) {
cli := r.getOrCreateLeaderClient(leaderAddr)

ctx, cancel := context.WithTimeout(r.handlerContext(), redisDispatchTimeout)
Expand Down Expand Up @@ -81,6 +93,20 @@ func (r *RedisServer) resolveLeaderRedisAddr(conn redcon.Conn) (string, bool) {
return leaderAddr, true
}

func (r *RedisServer) resolveLeaderRedisAddrForKey(conn redcon.Conn, key []byte) (string, bool) {
leader := r.coordinator.RaftLeaderForKey(key)
if leader == "" {
writeRedisError(conn, ErrLeaderNotFound)
return "", false
}
leaderAddr, ok := r.leaderRedis[leader]
if !ok || leaderAddr == "" {
conn.WriteError(fmt.Sprintf("ERR leader redis address unknown for raft address %s", leader))
return "", false
}
return leaderAddr, true
}

// execTxPipeline sends queue as a single TxPipelined batch and returns the
// per-command result handles together with any pipeline-level error.
func (r *RedisServer) execTxPipeline(ctx context.Context, cli *redis.Client, queue []redcon.Command) ([]*redis.Cmd, error) {
Expand Down Expand Up @@ -115,13 +141,7 @@ func handleProxyTxnError(conn redcon.Conn, err error) bool {
conn.WriteError(errRedisHeavyCommandPoolFull.Error())
return true
}
var netErr net.Error
if isTransientLeaderRedisError(err) ||
errors.Is(err, context.DeadlineExceeded) ||
errors.Is(err, context.Canceled) ||
errors.Is(err, io.EOF) ||
errors.Is(err, io.ErrUnexpectedEOF) ||
errors.As(err, &netErr) {
if isTerminalProxyTxnError(err) {
writeRedisError(conn, err)
return true
}
Expand All @@ -143,6 +163,10 @@ func handleProxyTxnCommandError(conn redcon.Conn, cmds []*redis.Cmd) bool {
conn.WriteError(errRedisHeavyCommandPoolFull.Error())
return true
}
if isRedisExecTerminalProxyError(err) {
writeRedisError(conn, err)
return true
}
if isTransientLeaderRedisError(err) {
writeRedisError(conn, err)
return true
Expand All @@ -151,6 +175,23 @@ func handleProxyTxnCommandError(conn redcon.Conn, cmds []*redis.Cmd) bool {
return false
}

func isTerminalProxyTxnError(err error) bool {
var netErr net.Error
return isRedisExecTerminalProxyError(err) ||
isTransientLeaderRedisError(err) ||
errors.Is(err, context.DeadlineExceeded) ||
errors.Is(err, context.Canceled) ||
errors.Is(err, io.EOF) ||
errors.Is(err, io.ErrUnexpectedEOF) ||
errors.As(err, &netErr)
}

func isRedisExecTerminalProxyError(err error) bool {
return err != nil &&
(errors.Is(err, errRedisExecRouteChangedAfterAmbiguousAttempt) ||
err.Error() == errRedisExecRouteChangedAfterAmbiguousAttempt.Error())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Promote split-leader EXEC errors across proxies

If leadership changes after one node chooses a common remote Redis target but before that target resolves its own read-fence groups, the target can return errRedisExecSplitShardLeaders as a top-level EXEC error. go-redis surfaces that decoded error on the pipeline and queued commands, but this predicate recognizes only the ambiguous-route-change text, so the proxy falls through to writeProxyCmdsResult and turns a transaction-level failure into an EXEC result array of command errors. Match the exact split-leader error here as another terminal proxy error.

Useful? React with 👍 / 👎.

}

// writeProxyCmdsResult writes an EXEC-style array reply for the given pipeline
// command handles. For any other non-nil per-command errors, each cmd carries
// its own result, which is the correct Redis EXEC semantics.
Expand All @@ -161,8 +202,8 @@ func writeProxyCmdsResult(conn redcon.Conn, cmds []*redis.Cmd) {
}
}

func (r *RedisServer) proxyLRange(key []byte, startRaw, endRaw []byte) ([]string, error) {
leader := r.coordinator.RaftLeaderForKey(key)
func (r *RedisServer) proxyLRange(key, routingKey []byte, startRaw, endRaw []byte) ([]string, error) {
leader := r.coordinator.RaftLeaderForKey(routingKey)
if leader == "" {
return nil, ErrLeaderNotFound
}
Expand Down
59 changes: 49 additions & 10 deletions adapter/redis_retry.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,11 +44,25 @@ var (
)

func isRetryableRedisTxnErr(err error) bool {
return isReusableRedisTxnErr(err) || isRedisComposedRouteErr(err)
}

func isReusableRedisTxnErr(err error) bool {
return errors.Is(err, store.ErrWriteConflict) ||
errors.Is(err, kv.ErrTxnLocked) ||
wireRedisTxnErrKind(err) == redisTxnWireErrLocked
}

func isRedisComposedRouteErr(err error) bool {
if errors.Is(err, kv.ErrComposed1Violation) ||
errors.Is(err, kv.ErrComposed1VersionGCd) {
return true
}
parsed := parseWireRedisTxnErr(err)
return errors.Is(parsed, kv.ErrComposed1Violation) ||
errors.Is(parsed, kv.ErrComposed1VersionGCd)
}

func retryPolicyForRedisTxnErr(err error) redisTxnRetryPolicy {
if errors.Is(err, kv.ErrTxnLocked) || wireRedisTxnErrKind(err) == redisTxnWireErrLocked {
return redisTxnLockedRetryPolicy
Expand All @@ -66,13 +80,13 @@ const (

// parseWireRedisTxnErr restores transaction error typing after an internal
// leader redirect crosses gRPC. Forward currently returns transaction failures
// as a status, which strips the typed ErrWriteConflict / ErrTxnLocked chain.
// Match only the exact server-generated key error envelope and normalize the
// storage key before rebuilding the typed error. Wire write conflicts are not
// generally retryable: a lost forwarding response can turn an already-applied
// write into a later self-conflict. Reuse-aware callers explicitly normalize
// them before retrying; all other callers return the normalized error without
// replaying the operation or exposing the internal key layout.
// as a status, which strips the typed ErrWriteConflict / ErrTxnLocked /
// Composed-1 sentinel chain. Match only known server-generated envelopes.
// Wire write conflicts are not generally retryable: a lost forwarding response
// can turn an already-applied write into a later self-conflict. Reuse-aware
// callers explicitly normalize them before retrying; all other callers return
// the normalized error without replaying the operation or exposing the internal
// key layout.
func wireRedisTxnStatus(err error) (*status.Status, bool) {
type grpcStatusCarrier interface {
GRPCStatus() *status.Status
Expand All @@ -94,6 +108,31 @@ func parseWireRedisTxnErr(err error) error {
return nil
}
msg := st.Message()
if parsed := parseWireRedisTxnKeyErr(msg); parsed != nil {
return parsed
}
if isWireComposedRouteMessage(msg, kv.ErrComposed1Violation) {
return errors.WithStack(kv.ErrComposed1Violation)
}
if isWireComposedRouteMessage(msg, kv.ErrComposed1VersionGCd) {
return errors.WithStack(kv.ErrComposed1VersionGCd)
}
return nil
}

func isWireComposedRouteMessage(msg string, sentinel error) bool {
sentinelMsg := sentinel.Error()
if msg == sentinelMsg {
return true
}
if !strings.HasSuffix(msg, ": "+sentinelMsg) {
return false
}
return strings.HasPrefix(msg, "observed-version v=") ||
strings.HasPrefix(msg, "current-version v=")
}

func parseWireRedisTxnKeyErr(msg string) error {
if !strings.HasPrefix(msg, "key: ") {
return nil
}
Expand Down Expand Up @@ -274,15 +313,15 @@ func normalizeRetryableRedisTxnKey(key []byte) []byte {
if userKey := redisTxnWideFenceUserKey(key); userKey != nil {
return userKey
}
if store.IsListMetaKey(key) || store.IsListItemKey(key) {
return store.ExtractListUserKey(key)
}
if store.IsListMetaDeltaKey(key) {
return store.ExtractListUserKeyFromDelta(key)
}
if store.IsListClaimKey(key) {
return store.ExtractListUserKeyFromClaim(key)
}
if store.IsListMetaKey(key) || store.IsListItemKey(key) {
return store.ExtractListUserKey(key)
}
if wideKey, ok := normalizeWideColumnKey(key); ok {
return wideKey
}
Expand Down
Loading
Loading