diff --git a/adapter/redis_peer_limiter_test.go b/adapter/redis_peer_limiter_test.go index 2f2267889..ddd3b608e 100644 --- a/adapter/redis_peer_limiter_test.go +++ b/adapter/redis_peer_limiter_test.go @@ -86,7 +86,7 @@ func TestRedisLeaderClientPoolStaysBelowPeerLimit(t *testing.T) { } func TestRedisLeaderClientPoolUsesSmallDefault(t *testing.T) { - server := NewRedisServer(nil, "", nil, nil, nil, nil, WithRedisPerPeerConnectionLimit(8)) + server := NewRedisServer(nil, "", nil, nil, nil, nil, WithRedisPerPeerConnectionLimit(defaultRedisPerPeerConnectionCap)) client := server.getOrCreateLeaderClient("127.0.0.1:6379") defer client.Close() @@ -103,8 +103,8 @@ func TestRedisLeaderClientPoolsSharePeerBudget(t *testing.T) { }{ {name: "low cap", limit: 2, wantNormal: 1, wantBlocking: 1}, {name: "four cap", limit: 4, wantNormal: 2, wantBlocking: 2}, - {name: "default cap", limit: 8, wantNormal: 4, wantBlocking: 4}, - {name: "raised cap", limit: defaultRedisPerPeerConnectionCap, wantNormal: 4, wantBlocking: 4}, + {name: "legacy cap", limit: 8, wantNormal: defaultRedisLeaderClientPoolSize, wantBlocking: defaultRedisBlockingLeaderClientPoolSize}, + {name: "default cap", limit: defaultRedisPerPeerConnectionCap, wantNormal: defaultRedisLeaderClientPoolSize, wantBlocking: defaultRedisBlockingLeaderClientPoolSize}, } { t.Run(tc.name, func(t *testing.T) { server := NewRedisServer(nil, "", nil, nil, nil, nil, WithRedisPerPeerConnectionLimit(tc.limit)) @@ -116,7 +116,7 @@ func TestRedisLeaderClientPoolsSharePeerBudget(t *testing.T) { } func TestRedisBlockingLeaderClientUsesDedicatedBudgetedPool(t *testing.T) { - server := NewRedisServer(nil, "", nil, nil, nil, nil, WithRedisPerPeerConnectionLimit(8)) + server := NewRedisServer(nil, "", nil, nil, nil, nil, WithRedisPerPeerConnectionLimit(defaultRedisPerPeerConnectionCap)) shared := server.getOrCreateLeaderClient("127.0.0.1:6379") defer shared.Close() blocking := server.getOrCreateBlockingLeaderClient("127.0.0.1:6379") diff --git a/docs/design/2026_04_24_implemented_workload_isolation.md b/docs/design/2026_04_24_implemented_workload_isolation.md index a929b0c2d..8ee69b3ee 100644 --- a/docs/design/2026_04_24_implemented_workload_isolation.md +++ b/docs/design/2026_04_24_implemented_workload_isolation.md @@ -26,9 +26,11 @@ Implementation status: `adapter/redis_peer_limiter.go`, wired through `RedisServer.Run` accept and close hooks. Default cap is 8 per peer IP and is configurable via `ELASTICKV_REDIS_PER_PEER_CONNECTIONS` / - `WithRedisPerPeerConnectionLimit`. Redis leader-proxy clients use a small - explicit go-redis pool below that default cap, and Pub/Sub detached sockets - stay counted until their Pub/Sub cleanup path closes. + `WithRedisPerPeerConnectionLimit`. The redis-proxy deployment can raise the + cap explicitly on ElasticKV nodes before increasing the proxy's ElasticKV + pool size. Redis leader-proxy clients use a small explicit go-redis pool + below that default cap, and Pub/Sub detached sockets stay counted until their + Pub/Sub cleanup path closes. - Shipped: Layer 4 stream entry-per-key layout in `store/stream_helpers.go`, `adapter/redis_stream_cmds.go`, and `adapter/redis_compat_helpers.go`. XREAD now range-scans `!stream|entry|...` after the requested ID instead of diff --git a/docs/redis-proxy-deployment.md b/docs/redis-proxy-deployment.md index c613d9c1c..f9027487b 100644 --- a/docs/redis-proxy-deployment.md +++ b/docs/redis-proxy-deployment.md @@ -414,7 +414,7 @@ groups: | Parameter | Value | Description | |-----------|-------|-------------| | Redis connection pool size | 128 | Default go-redis pool size for Redis | -| ElasticKV connection pool size | 64 | Default per-leader command pool; leave server per-peer headroom for dedicated PubSub connections | +| ElasticKV connection pool size | 64 | Default per-leader command pool; regular and blocking replay pools share the server per-peer budget, while PubSub uses dedicated sockets outside the pool | | Dial timeout | 5s | Backend connection timeout | | Read timeout | 3s | Backend read timeout | | Write timeout | 3s | Backend write timeout | @@ -439,7 +439,7 @@ Recommended shutdown order: `redis-proxy -> application -> Redis / ElasticKV`. ### Secondary writes are falling behind - Check `proxy_async_queue_depth`, `proxy_async_queue_delay_seconds`, and `proxy_async_drops_by_queue_total` before increasing concurrency. - Check `proxy_backend_pool_pending_requests` and the `waits`/`timeouts` pool events. Pool waits mean concurrency is too high for the configured pool. -- Keep `ELASTICKV_REDIS_PER_PEER_CONNECTIONS` above `-elastickv-pool-size`; PubSub and shadow PubSub use dedicated connections outside the command pool. Keep `-secondary-write-concurrency` at or below the pool size. +- Keep `ELASTICKV_REDIS_PER_PEER_CONNECTIONS` above `-elastickv-pool-size`; PubSub and shadow PubSub use dedicated connections outside the command pool and detached PubSub sockets may remain counted until cleanup. Keep `-secondary-write-concurrency` at or below the pool size. - A sustained `expired` rate means secondary throughput is below ingress. Increasing queue size only delays the loss; profile ElasticKV before raising concurrency. ### High divergence count diff --git a/proxy/backend.go b/proxy/backend.go index 599e995dc..704cfcc45 100644 --- a/proxy/backend.go +++ b/proxy/backend.go @@ -72,11 +72,12 @@ func DefaultBackendOptions() BackendOptions { } // DefaultElasticKVBackendOptions returns defaults for proxy backends that -// connect to ElasticKV's Redis adapter. Production dual-write deployments -// should run the cluster with ELASTICKV_REDIS_PER_PEER_CONNECTIONS above this -// pool size because PubSub and shadow PubSub use dedicated connections outside -// the go-redis command pool. Lower the proxy pool instead for clusters that -// keep the server-side per-peer cap below the default. +// connect to ElasticKV's Redis adapter. The regular command pool and blocking +// replay pool share ElasticKV's per-peer connection budget. PubSub and shadow +// PubSub use dedicated sockets outside those pools, and detached sockets may +// remain counted until cleanup. Run the cluster with +// ELASTICKV_REDIS_PER_PEER_CONNECTIONS above this pool size before using the +// default in production, or lower the proxy pool to fit the configured cap. func DefaultElasticKVBackendOptions() BackendOptions { opts := DefaultBackendOptions() opts.PoolSize = defaultElasticKVPoolSize diff --git a/proxy/leader_aware_backend.go b/proxy/leader_aware_backend.go index 05bd6e6b0..0b83aebc2 100644 --- a/proxy/leader_aware_backend.go +++ b/proxy/leader_aware_backend.go @@ -59,6 +59,7 @@ type LeaderAwareRedisBackend struct { refreshMu sync.Mutex refreshDone chan struct{} refreshClosed bool + lastRefreshAt time.Time refreshCtx context.Context refreshCancel context.CancelFunc @@ -75,8 +76,8 @@ type LeaderAwareRedisBackend struct { } // NewLeaderAwareRedisBackend creates a LeaderAwareRedisBackend with the given -// seed addresses. The first seed is used as the initial target until the -// first refresh completes. At least one seed is required. +// seed addresses. The first command waits for leader discovery instead of +// sending traffic to a seed that may be down. At least one seed is required. func NewLeaderAwareRedisBackend(seeds []string, name string, opts BackendOptions, logger *slog.Logger) *LeaderAwareRedisBackend { return NewLeaderAwareRedisBackendWithInterval(seeds, name, opts, defaultLeaderRefreshInterval, defaultLeaderRefreshTimeout, logger) } @@ -106,7 +107,7 @@ func NewLeaderAwareRedisBackendWithInterval(seeds []string, name string, opts Ba clients: make(map[string]*redis.Client, len(normalized)), clientOrder: make([]string, 0, len(normalized)), seedProtect: seedProtect, - leader: normalized[0], + leader: "", stopCh: make(chan struct{}), done: make(chan struct{}), refreshCh: make(chan struct{}, 1), @@ -209,6 +210,7 @@ func (b *LeaderAwareRedisBackend) runLeaderRefresh(done chan struct{}) { b.refreshMu.Lock() b.refreshDone = nil + b.lastRefreshAt = time.Now() close(done) b.refreshMu.Unlock() } @@ -244,6 +246,30 @@ func (b *LeaderAwareRedisBackend) RefreshLeaderNow(ctx context.Context) { b.refreshLeader(ctx) } +func (b *LeaderAwareRedisBackend) refreshLeaderNowIfDue(ctx context.Context) { + b.refreshMu.Lock() + if b.refreshClosed { + b.refreshMu.Unlock() + return + } + done := b.refreshDone + if done == nil { + if !b.lastRefreshAt.IsZero() && time.Since(b.lastRefreshAt) < b.refreshInterval { + b.refreshMu.Unlock() + return + } + done = make(chan struct{}) + b.refreshDone = done + go b.runLeaderRefresh(done) + } + b.refreshMu.Unlock() + + select { + case <-done: + case <-ctx.Done(): + } +} + func (b *LeaderAwareRedisBackend) probeLeader(ctx context.Context, addr string) (string, error) { cli := b.getOrCreateClient(addr) if cli == nil { @@ -389,6 +415,29 @@ func (b *LeaderAwareRedisBackend) currentClient() *redis.Client { return b.clients[b.leader] } +func (b *LeaderAwareRedisBackend) currentClientOrRefresh(ctx context.Context) *redis.Client { + cli := b.currentClient() + if cli != nil { + return cli + } + b.refreshLeaderNowIfDue(ctx) + return b.currentClient() +} + +func (b *LeaderAwareRedisBackend) firstSeedClient() *redis.Client { + b.mu.RLock() + defer b.mu.RUnlock() + if b.closed { + return nil + } + for _, seed := range b.seeds { + if cli := b.clients[seed]; cli != nil { + return cli + } + } + return nil +} + // Do forwards a single command to the current leader. NOTLEADER refreshes the // cached leader for the next command, but the current command is not replayed: // leadership-loss errors can be returned after an operation has already applied. @@ -404,7 +453,7 @@ func (b *LeaderAwareRedisBackend) Do(ctx context.Context, args ...any) *redis.Cm } func (b *LeaderAwareRedisBackend) doOnce(ctx context.Context, args ...any) *redis.Cmd { - cli := b.currentClient() + cli := b.currentClientOrRefresh(ctx) if cli == nil { cmd := redis.NewCmd(ctx, args...) cmd.SetErr(ErrNoLeaderBackend) @@ -426,7 +475,7 @@ func (b *LeaderAwareRedisBackend) DoWithTimeout(ctx context.Context, timeout tim } func (b *LeaderAwareRedisBackend) doWithTimeoutOnce(ctx context.Context, timeout time.Duration, args ...any) *redis.Cmd { - cli := b.currentClient() + cli := b.currentClientOrRefresh(ctx) if cli == nil { cmd := redis.NewCmd(ctx, args...) cmd.SetErr(ErrNoLeaderBackend) @@ -437,7 +486,7 @@ func (b *LeaderAwareRedisBackend) doWithTimeoutOnce(ctx context.Context, timeout // Pipeline forwards a batch to the current leader. func (b *LeaderAwareRedisBackend) Pipeline(ctx context.Context, cmds [][]any) ([]*redis.Cmd, error) { - cli := b.currentClient() + cli := b.currentClientOrRefresh(ctx) if cli == nil { return nil, ErrNoLeaderBackend } @@ -517,7 +566,10 @@ func isLeaderRefreshTransportError(err error) bool { // NewPubSub opens a subscribe connection on the current leader. func (b *LeaderAwareRedisBackend) NewPubSub(ctx context.Context) *redis.PubSub { - cli := b.currentClient() + cli := b.currentClientOrRefresh(ctx) + if cli == nil { + cli = b.firstSeedClient() + } if cli == nil { return nil } diff --git a/proxy/leader_aware_backend_test.go b/proxy/leader_aware_backend_test.go index c42cfbeed..5d40cefd7 100644 --- a/proxy/leader_aware_backend_test.go +++ b/proxy/leader_aware_backend_test.go @@ -150,13 +150,14 @@ func (n *fakeElasticKVNode) handleConn(conn net.Conn) { ) _, _ = fmt.Fprintf(conn, "$%d\r\n%s\r\n", len(body), body) default: - n.writeCommandResponse(conn, cmd) + n.writeCommandResponse(conn, args) } } } -func (n *fakeElasticKVNode) writeCommandResponse(conn net.Conn, cmd string) { +func (n *fakeElasticKVNode) writeCommandResponse(conn net.Conn, args []string) { n.commands.Add(1) + cmd := strings.ToUpper(args[0]) if leader := n.Leader(); leader != "" && leader != n.addr { _, _ = conn.Write([]byte("-NOTLEADER etcd raft engine is not leader\r\n")) return @@ -165,9 +166,32 @@ func (n *fakeElasticKVNode) writeCommandResponse(conn net.Conn, cmd string) { _, _ = conn.Write([]byte("-NOTLEADER user script\r\n")) return } + switch cmd { + case "SUBSCRIBE": + n.writeSubscribeAck(conn, "subscribe", args[1:]) + return + case "PSUBSCRIBE": + n.writeSubscribeAck(conn, "psubscribe", args[1:]) + return + } _, _ = conn.Write([]byte("+OK\r\n")) } +func (n *fakeElasticKVNode) writeSubscribeAck(conn net.Conn, kind string, names []string) { + if len(names) == 0 { + _, _ = conn.Write([]byte("+OK\r\n")) + return + } + for i, name := range names { + _, _ = fmt.Fprintf(conn, + "*3\r\n$%d\r\n%s\r\n$%d\r\n%s\r\n:%d\r\n", + len(kind), kind, + len(name), name, + i+1, + ) + } +} + // readRESPArray reads a single RESP array of bulk strings from rd. func readRESPArray(rd *bufio.Reader) ([]string, error) { header, err := rd.ReadString('\n') @@ -553,6 +577,138 @@ func TestLeaderAwareRedisBackend_EvictsOldestNonProtectedClient(t *testing.T) { assert.Contains(t, addrs, seed.addr, "seed must never be evicted") } +func TestLeaderAwareRedisBackend_InitialCommandWaitsForLeaderDiscovery(t *testing.T) { + tests := []struct { + name string + run func(context.Context, *LeaderAwareRedisBackend) error + }{ + { + name: "Do", + run: func(ctx context.Context, backend *LeaderAwareRedisBackend) error { + return backend.Do(ctx, "SET", "k", "v").Err() + }, + }, + { + name: "Pipeline", + run: func(ctx context.Context, backend *LeaderAwareRedisBackend) error { + _, err := backend.Pipeline(ctx, [][]any{{"SET", "k", "v"}}) + return err + }, + }, + { + name: "NewPubSub", + run: func(ctx context.Context, backend *LeaderAwareRedisBackend) error { + ps := backend.NewPubSub(ctx) + if ps == nil { + return ErrNoLeaderBackend + } + defer ps.Close() + return ps.Subscribe(ctx, "events") + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + aliveLeader := newFakeElasticKVNode(t) + aliveLeader.SetLeader(aliveLeader.addr) + gate := make(chan struct{}) + aliveLeader.SetInfoGate(gate) + + backend := NewLeaderAwareRedisBackendWithInterval( + []string{"127.0.0.1:1", aliveLeader.addr}, + "elastickv", + DefaultBackendOptions(), + time.Hour, 200*time.Millisecond, + testLogger, + ) + t.Cleanup(func() { _ = backend.Close() }) + + require.Empty(t, backend.CurrentLeader(), "initial commands must not target the first seed before discovery") + + done := make(chan error, 1) + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + done <- tt.run(ctx, backend) + }() + + select { + case err := <-done: + t.Fatalf("operation returned before leader discovery completed: %v", err) + case <-time.After(50 * time.Millisecond): + } + require.Equal(t, int64(0), aliveLeader.commands.Load(), "operation must not run before discovery completes") + + close(gate) + var err error + require.Eventually(t, func() bool { + select { + case err = <-done: + return true + default: + return false + } + }, 2*time.Second, 10*time.Millisecond) + require.NoError(t, err) + require.Equal(t, int64(1), aliveLeader.commands.Load(), "initial operation must reach the discovered leader") + require.Equal(t, aliveLeader.addr, backend.CurrentLeader()) + }) + } +} + +func TestLeaderAwareRedisBackend_NewPubSubFallsBackToSeedWhenNoLeaderDiscovered(t *testing.T) { + seed := newFakeElasticKVNode(t) + + backend := NewLeaderAwareRedisBackendWithInterval( + []string{seed.addr}, + "elastickv", + DefaultBackendOptions(), + time.Hour, 50*time.Millisecond, + testLogger, + ) + t.Cleanup(func() { _ = backend.Close() }) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + ps := backend.NewPubSub(ctx) + require.NotNil(t, ps) + t.Cleanup(func() { _ = ps.Close() }) + + require.NoError(t, ps.Subscribe(ctx, "events")) + require.Eventually(t, func() bool { + return seed.commands.Load() == 1 + }, time.Second, 10*time.Millisecond, "PubSub fallback should use the seed instead of returning nil") + require.Empty(t, backend.CurrentLeader()) +} + +func TestLeaderAwareRedisBackend_UnknownLeaderDoesNotProbeOnEveryCommand(t *testing.T) { + seed := newFakeElasticKVNode(t) + + backend := NewLeaderAwareRedisBackendWithInterval( + []string{seed.addr}, + "elastickv", + DefaultBackendOptions(), + time.Hour, 50*time.Millisecond, + testLogger, + ) + t.Cleanup(func() { _ = backend.Close() }) + + require.Eventually(t, func() bool { + return seed.infoCalls.Load() > 0 + }, time.Second, 10*time.Millisecond, "initial background probe should run") + require.Empty(t, backend.CurrentLeader()) + before := seed.infoCalls.Load() + + for i := 0; i < 3; i++ { + err := backend.Do(context.Background(), "SET", "k", "v").Err() + require.ErrorIs(t, err, ErrNoLeaderBackend) + } + + require.Equal(t, before, seed.infoCalls.Load(), "commands must not start a fresh probe inside the refresh interval") + require.Equal(t, int64(0), seed.commands.Load(), "commands must not fall back to a seed when leader discovery returned no leader") +} + func TestLeaderAwareRedisBackend_FallsBackToSeedOnProbeFailure(t *testing.T) { // Seed 1 is unreachable; seed 2 is alive and reports itself as leader. aliveLeader := newFakeElasticKVNode(t) diff --git a/proxy/proxy_test.go b/proxy/proxy_test.go index 01389e6d1..22936507f 100644 --- a/proxy/proxy_test.go +++ b/proxy/proxy_test.go @@ -813,7 +813,7 @@ func TestDualWriter_Blocking_BZPopReplayShortTimeoutStillAttemptsZRem(t *testing func TestNoEffectReplayRetryLimitIncludesJitterBudget(t *testing.T) { limit := noEffectReplayRetryLimit(context.Background(), blockingReplayNoEffectRetryWindow) - assert.Positive(t, limit) + assert.Equal(t, 5, limit) var spent time.Duration backoff := compactedRetryInitialBackoff