//go:build integration package redisprovider import ( "context" "fmt" "os" "sync/atomic" "testing" "time" "github.com/redis/go-redis/v9" controllerProvider "proxy-pool/internal/controller/provider" ) var integrationNamespaceSequence atomic.Uint64 func TestRedisCoordinatorElectsOneLeaderAndFencesFailover(t *testing.T) { fixture := newRedisFixture(t) first := fixture.coordinator(t, "controller-a") second := fixture.coordinator(t, "controller-b") limits := controllerProvider.CoordinationLimits{ RequestInterval: 50 * time.Millisecond, MaxInFlight: 1, MaxAttemptDuration: 300 * time.Millisecond, } ctxA, cancelA := context.WithCancel(context.Background()) ctxB, cancelB := context.WithCancel(context.Background()) defer cancelA() defer cancelB() started := make(chan leadershipFixture, 4) var active atomic.Int64 var maximum atomic.Int64 work := func(holder string) func(context.Context, controllerProvider.LeaderSession) error { return func(ctx context.Context, session controllerProvider.LeaderSession) error { current := active.Add(1) for { observed := maximum.Load() if current <= observed || maximum.CompareAndSwap(observed, current) { break } } started <- leadershipFixture{holder: holder, fence: session.Fence()} <-ctx.Done() active.Add(-1) return nil } } doneA := make(chan error, 1) doneB := make(chan error, 1) go func() { doneA <- first.RunLeader(ctxA, "provider-a", limits, work("controller-a")) }() go func() { doneB <- second.RunLeader(ctxB, "provider-a", limits, work("controller-b")) }() initial := receiveLeadership(t, started) time.Sleep(200 * time.Millisecond) select { case duplicate := <-started: t.Fatalf("simultaneous leaders started: first=%+v duplicate=%+v", initial, duplicate) default: } if got := maximum.Load(); got != 1 { t.Fatalf("maximum simultaneous leaders = %d, want 1", got) } if initial.holder == "controller-a" { cancelA() } else { cancelB() } replacement := receiveLeadership(t, started) if replacement.holder == initial.holder { t.Fatalf("replacement holder = %q, want the other controller", replacement.holder) } if replacement.fence.Generation != initial.fence.Generation || replacement.fence.Epoch <= initial.fence.Epoch { t.Fatalf("replacement fence = %+v, initial = %+v", replacement.fence, initial.fence) } if got := maximum.Load(); got != 1 { t.Fatalf("maximum simultaneous leaders after failover = %d, want 1", got) } cancelA() cancelB() waitRunner(t, doneA) waitRunner(t, doneB) } func TestRedisLeaderSessionEnforcesGlobalIntervalAndInFlightLimit(t *testing.T) { fixture := newRedisFixture(t) coordinator := fixture.coordinator(t, "controller-a") limits := controllerProvider.CoordinationLimits{ RequestInterval: 250 * time.Millisecond, MaxInFlight: 1, MaxAttemptDuration: time.Second, } ctx, cancel := context.WithCancel(context.Background()) defer cancel() sessions := make(chan controllerProvider.LeaderSession, 1) done := make(chan error, 1) go func() { done <- coordinator.RunLeader(ctx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { sessions <- session <-workCtx.Done() return nil }) }() session := receiveSession(t, sessions) first, available, err := session.AcquireFetch(context.Background(), 1) if err != nil || !available { t.Fatalf("first AcquireFetch() = (%v, %t, %v)", first, available, err) } startedAt := time.Now() secondResult := make(chan permitResultFixture, 1) go func() { permit, permitAvailable, acquireErr := session.AcquireFetch(context.Background(), 1) secondResult <- permitResultFixture{permit: permit, available: permitAvailable, err: acquireErr} }() select { case result := <-secondResult: t.Fatalf("second AcquireFetch() returned before release: %+v", result) case <-time.After(100 * time.Millisecond): } if err := first.Complete(context.Background(), 1); err != nil { t.Fatalf("first Complete(): %v", err) } result := receivePermit(t, secondResult) if result.err != nil || !result.available || result.permit == nil { t.Fatalf("second AcquireFetch() = (%v, %t, %v)", result.permit, result.available, result.err) } if elapsed := time.Since(startedAt); elapsed < 200*time.Millisecond { t.Fatalf("global request interval = %s, want at least 200ms", elapsed) } if err := result.permit.Cancel(context.Background()); err != nil { t.Fatalf("second Cancel(): %v", err) } if err := result.permit.Cancel(context.Background()); err != nil { t.Fatalf("idempotent second Cancel(): %v", err) } cancel() waitRunner(t, done) } func TestRedisLeaderSessionPreservesFetchQuotaAcrossFailover(t *testing.T) { fixture := newRedisFixture(t) limits := controllerProvider.CoordinationLimits{ MaxInFlight: 1, MaxAttemptDuration: time.Second, MaxTotal: 2, } firstCtx, cancelFirst := context.WithCancel(context.Background()) firstSessions := make(chan controllerProvider.LeaderSession, 1) firstDone := make(chan error, 1) go func() { firstDone <- fixture.coordinator(t, "controller-a").RunLeader( firstCtx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { firstSessions <- session <-workCtx.Done() return nil }, ) }() firstSession := receiveSession(t, firstSessions) firstPermit, available, err := firstSession.AcquireFetch(context.Background(), 2) if err != nil || !available || firstPermit == nil { t.Fatalf("first AcquireFetch() = (%v, %t, %v)", firstPermit, available, err) } firstFence := firstSession.Fence() cancelFirst() waitRunner(t, firstDone) if err := firstPermit.Complete(context.Background(), 1); err != nil { t.Fatalf("Complete() after leadership loss: %v", err) } secondCtx, cancelSecond := context.WithCancel(context.Background()) defer cancelSecond() secondSessions := make(chan controllerProvider.LeaderSession, 1) secondDone := make(chan error, 1) go func() { secondDone <- fixture.coordinator(t, "controller-b").RunLeader( secondCtx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { secondSessions <- session <-workCtx.Done() return nil }, ) }() secondSession := receiveSession(t, secondSessions) secondFence := secondSession.Fence() if secondFence.Generation != firstFence.Generation || secondFence.Epoch <= firstFence.Epoch { t.Fatalf("second fence = %+v, first = %+v", secondFence, firstFence) } secondPermit, available, err := secondSession.AcquireFetch(context.Background(), 1) if err != nil || !available || secondPermit == nil { t.Fatalf("second AcquireFetch() = (%v, %t, %v)", secondPermit, available, err) } if err := secondPermit.Complete(context.Background(), 1); err != nil { t.Fatalf("second Complete(): %v", err) } exhaustedPermit, available, err := secondSession.AcquireFetch(context.Background(), 1) if err != nil || available || exhaustedPermit != nil { t.Fatalf("exhausted AcquireFetch() = (%v, %t, %v), want unavailable", exhaustedPermit, available, err) } cancelSecond() waitRunner(t, secondDone) } func TestRedisFetchQuotaSettlementIsIdempotentAndCancellationRefundsReservation(t *testing.T) { fixture := newRedisFixture(t) limits := controllerProvider.CoordinationLimits{ MaxInFlight: 2, MaxAttemptDuration: time.Second, MaxTotal: 2, } ctx, cancel := context.WithCancel(context.Background()) defer cancel() sessions := make(chan controllerProvider.LeaderSession, 1) done := make(chan error, 1) go func() { done <- fixture.coordinator(t, "controller-a").RunLeader( ctx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { sessions <- session <-workCtx.Done() return nil }, ) }() session := receiveSession(t, sessions) cancelled, available, err := session.AcquireFetch(context.Background(), 2) if err != nil || !available || cancelled == nil { t.Fatalf("cancelled AcquireFetch() = (%v, %t, %v)", cancelled, available, err) } if err := cancelled.Cancel(context.Background()); err != nil { t.Fatalf("Cancel(): %v", err) } if err := cancelled.Cancel(context.Background()); err != nil { t.Fatalf("idempotent Cancel(): %v", err) } completed, available, err := session.AcquireFetch(context.Background(), 2) if err != nil || !available || completed == nil { t.Fatalf("completed AcquireFetch() = (%v, %t, %v)", completed, available, err) } if err := completed.Complete(context.Background(), 1); err != nil { t.Fatalf("Complete(): %v", err) } if err := completed.Complete(context.Background(), 1); err != nil { t.Fatalf("idempotent Complete(): %v", err) } if err := completed.Cancel(context.Background()); err != nil { t.Fatalf("Cancel() after Complete(): %v", err) } last, available, err := session.AcquireFetch(context.Background(), 1) if err != nil || !available || last == nil { t.Fatalf("last AcquireFetch() = (%v, %t, %v)", last, available, err) } if err := last.Complete(context.Background(), 1); err != nil { t.Fatalf("last Complete(): %v", err) } exhausted, available, err := session.AcquireFetch(context.Background(), 1) if err != nil || available || exhausted != nil { t.Fatalf("exhausted AcquireFetch() = (%v, %t, %v)", exhausted, available, err) } cancel() waitRunner(t, done) } func TestRedisExpiredFetchReservationIsConservativelyCharged(t *testing.T) { fixture := newRedisFixture(t) coordinator, err := New(fixture.client, Options{ Namespace: fixture.namespace, HolderID: "controller-a", LeaseTTL: 600 * time.Millisecond, RenewEvery: 150 * time.Millisecond, RetryInterval: 10 * time.Millisecond, PermitGrace: 10 * time.Millisecond, }) if err != nil { t.Fatalf("New(): %v", err) } limits := controllerProvider.CoordinationLimits{ MaxInFlight: 1, MaxAttemptDuration: 40 * time.Millisecond, MaxTotal: 1, } ctx, cancel := context.WithCancel(context.Background()) defer cancel() sessions := make(chan controllerProvider.LeaderSession, 1) done := make(chan error, 1) go func() { done <- coordinator.RunLeader(ctx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { sessions <- session <-workCtx.Done() return nil }) }() session := receiveSession(t, sessions) abandoned, available, err := session.AcquireFetch(context.Background(), 1) if err != nil || !available || abandoned == nil { t.Fatalf("abandoned AcquireFetch() = (%v, %t, %v)", abandoned, available, err) } time.Sleep(80 * time.Millisecond) exhausted, available, err := session.AcquireFetch(context.Background(), 1) if err != nil || available || exhausted != nil { t.Fatalf("post-expiry AcquireFetch() = (%v, %t, %v), want charged quota", exhausted, available, err) } cancel() waitRunner(t, done) } func TestRedisCoordinatorRebuildsWithNewGenerationAfterStateLoss(t *testing.T) { fixture := newRedisFixture(t) coordinator := fixture.coordinator(t, "controller-a") limits := controllerProvider.CoordinationLimits{ RequestInterval: 50 * time.Millisecond, MaxInFlight: 1, MaxAttemptDuration: time.Second, } ctx, cancel := context.WithCancel(context.Background()) defer cancel() started := make(chan controllerProvider.Fence, 4) done := make(chan error, 1) go func() { done <- coordinator.RunLeader(ctx, "provider-a", limits, func(workCtx context.Context, session controllerProvider.LeaderSession) error { started <- session.Fence() <-workCtx.Done() return nil }) }() first := receiveFence(t, started) fixture.deleteKeys(t) second := receiveFence(t, started) if second.Generation == first.Generation { t.Fatalf("generation after Redis state loss = %q, want a new generation", second.Generation) } cancel() waitRunner(t, done) } type redisFixture struct { client *redis.Client namespace string } func newRedisFixture(t *testing.T) redisFixture { t.Helper() redisURL := os.Getenv("PROXY_POOL_TEST_REDIS_URL") if redisURL == "" { t.Skip("PROXY_POOL_TEST_REDIS_URL is not set") } options, err := redis.ParseURL(redisURL) if err != nil { t.Fatalf("parse PROXY_POOL_TEST_REDIS_URL: %v", err) } client := redis.NewClient(options) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := client.Ping(ctx).Err(); err != nil { _ = client.Close() t.Fatalf("ping Redis: %v", err) } namespace := fmt.Sprintf("provider-it-%d-%d-%d", os.Getpid(), time.Now().UnixNano(), integrationNamespaceSequence.Add(1)) t.Cleanup(func() { redisFixture{client: client, namespace: namespace}.deleteKeys(t) _ = client.Close() }) return redisFixture{client: client, namespace: namespace} } func (fixture redisFixture) deleteKeys(t *testing.T) { t.Helper() cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 5*time.Second) defer cleanupCancel() var cursor uint64 for { keys, next, err := fixture.client.Scan(cleanupCtx, cursor, "pp:"+fixture.namespace+":*", 128).Result() if err != nil { t.Fatalf("scan Redis provider keys: %v", err) } if len(keys) > 0 { if err := fixture.client.Unlink(cleanupCtx, keys...).Err(); err != nil { t.Fatalf("remove Redis provider keys: %v", err) } } cursor = next if cursor == 0 { return } } } func (fixture redisFixture) coordinator(t *testing.T, holder string) *Adapter { t.Helper() adapter, err := New(fixture.client, Options{ Namespace: fixture.namespace, HolderID: holder, LeaseTTL: 600 * time.Millisecond, RenewEvery: 150 * time.Millisecond, RetryInterval: 20 * time.Millisecond, PermitGrace: 100 * time.Millisecond, }) if err != nil { t.Fatalf("New(%s): %v", holder, err) } return adapter } func receiveLeadership(t *testing.T, values <-chan leadershipFixture) leadershipFixture { t.Helper() select { case value := <-values: return value case <-time.After(5 * time.Second): t.Fatal("timed out waiting for leadership") return leadershipFixture{} } } type leadershipFixture struct { holder string fence controllerProvider.Fence } func receiveSession(t *testing.T, values <-chan controllerProvider.LeaderSession) controllerProvider.LeaderSession { t.Helper() select { case value := <-values: return value case <-time.After(5 * time.Second): t.Fatal("timed out waiting for leader session") return nil } } func receiveFence(t *testing.T, values <-chan controllerProvider.Fence) controllerProvider.Fence { t.Helper() select { case value := <-values: return value case <-time.After(5 * time.Second): t.Fatal("timed out waiting for provider fence") return controllerProvider.Fence{} } } func receivePermit(t *testing.T, values <-chan permitResultFixture) permitResultFixture { t.Helper() select { case value := <-values: return value case <-time.After(5 * time.Second): t.Fatal("timed out waiting for request permit") return permitResultFixture{} } } type permitResultFixture struct { permit controllerProvider.RequestPermit available bool err error } func waitRunner(t *testing.T, done <-chan error) { t.Helper() select { case err := <-done: if err != nil && err != context.Canceled { t.Fatalf("RunLeader() error = %v", err) } case <-time.After(5 * time.Second): t.Fatal("timed out waiting for RunLeader shutdown") } }