diff --git a/logservice/eventstore/pebble.go b/logservice/eventstore/pebble.go index f5f48bf71b..96caabc3cb 100644 --- a/logservice/eventstore/pebble.go +++ b/logservice/eventstore/pebble.go @@ -42,7 +42,7 @@ func newPebbleOptions(dbNum int) *pebble.Options { MaxOpenFiles: maxOpenFilesPerDB, - MaxConcurrentCompactions: func() int { return 6 }, + MaxConcurrentCompactions: func() int { return 3 }, // Decrease compaction frequency L0CompactionThreshold: 20, diff --git a/logservice/logpuller/region_admission_controller_test.go b/logservice/logpuller/region_admission_controller_test.go index c1c393916b..2633e99681 100644 --- a/logservice/logpuller/region_admission_controller_test.go +++ b/logservice/logpuller/region_admission_controller_test.go @@ -99,11 +99,10 @@ func TestRegionAdmissionControllerNormalWindow(t *testing.T) { require.True(t, req2.abort()) } -func TestRegionAdmissionControllerLowLagUsesMaxWindow(t *testing.T) { +func TestRegionAdmissionControllerHighPriorityUsesMaxWindow(t *testing.T) { controller := newTestRegionAdmissionController(1, 2) currentTs := oracle.GoTimeToTS(time.Now()) slowCheckpointTs := oracle.GoTimeToTS(time.Now().Add(-time.Hour)) - lowLagCheckpointTs := oracle.GoTimeToTS(time.Now().Add(-time.Minute)) submitRegionForAdmission(t, controller, prepareRegionForAdmission(createTestRegionInfo(1, 1), slowCheckpointTs), @@ -114,9 +113,9 @@ func TestRegionAdmissionControllerLowLagUsesMaxWindow(t *testing.T) { submitRegionForAdmission(t, controller, prepareRegionForAdmission(createTestRegionInfo(1, 2), slowCheckpointTs), currentTs) - submitRegionForAdmission(t, controller, - prepareRegionForAdmission(createTestRegionInfo(1, 3), lowLagCheckpointTs), - currentTs) + highPriorityRegion := prepareRegionForAdmission(createTestRegionInfo(1, 3), slowCheckpointTs) + highPriorityRegion.scanPriority = cdcpb.ScanPriority_SCAN_PRIORITY_HIGH + submitRegionForAdmission(t, controller, highPriorityRegion, currentTs) req2, err := controller.pop(t.Context(), nil) require.NoError(t, err) @@ -152,8 +151,7 @@ func TestRegionAdmissionControllerPrioritizesHighPriorityRegion(t *testing.T) { currentTs) highPriorityRegion := prepareRegionForAdmission(createTestRegionInfo(1, 3), slowCheckpointTs) highPriorityRegion.scanPriority = cdcpb.ScanPriority_SCAN_PRIORITY_HIGH - submitRegionForAdmission(t, controller, - highPriorityRegion, currentTs) + submitRegionForAdmission(t, controller, highPriorityRegion, currentTs) req2, err := controller.pop(t.Context(), nil) require.NoError(t, err) diff --git a/logservice/logpuller/region_event_handler_test.go b/logservice/logpuller/region_event_handler_test.go index 38a2b50f2e..bbdd5cbee4 100644 --- a/logservice/logpuller/region_event_handler_test.go +++ b/logservice/logpuller/region_event_handler_test.go @@ -92,7 +92,7 @@ func TestHandleEventEntryEventOutOfOrder(t *testing.T) { false, ) region.lockedRangeState = ®ionlock.LockedRangeState{} - state := newRegionFeedState(region, 1, worker, nil) + state := newRegionFeedState(region, 1, worker, nil, nil) // Receive prewrite2 with empty value. { @@ -225,7 +225,7 @@ func TestHandleResolvedTs(t *testing.T) { worker := ®ionRequestWorker{ tracker: newRegionTracker(), } - state1 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(1, 1, 1)}, uint64(subID1), worker, nil) + state1 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(1, 1, 1)}, uint64(subID1), worker, nil, nil) { span := heartbeatpb.TableSpan{ TableID: 100, @@ -249,7 +249,7 @@ func TestHandleResolvedTs(t *testing.T) { } subID2 := SubscriptionID(2) - state2 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(2, 2, 2)}, uint64(subID2), worker, nil) + state2 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(2, 2, 2)}, uint64(subID2), worker, nil, nil) { span := heartbeatpb.TableSpan{ TableID: 100, @@ -273,7 +273,7 @@ func TestHandleResolvedTs(t *testing.T) { } subID3 := SubscriptionID(3) - state3 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(3, 3, 3)}, uint64(subID3), worker, nil) + state3 := newRegionFeedState(regionInfo{verID: tikv.NewRegionVerID(3, 3, 3)}, uint64(subID3), worker, nil, nil) { span := heartbeatpb.TableSpan{ TableID: 100, @@ -378,6 +378,7 @@ func TestHandleResolvedTsThrottled(t *testing.T) { 1, worker, nil, + nil, ) require.Equal(t, uint64(200), handleResolvedTs(span, state, 300)) @@ -430,6 +431,42 @@ func TestHandleEntriesReleasesMemoryAfterDownstreamCallback(t *testing.T) { require.Zero(t, quotaState.used) } +func TestRegionEventHandlerInitializedResetsRecoveryState(t *testing.T) { + span := &subscribedSpan{ + subID: 1, + span: heartbeatpb.TableSpan{TableID: 1}, + advanceResolvedTs: func(uint64) {}, + } + failureHandler := newRegionFailureHandler(nil, func(*subscribedSpan) {}, func(context.Context, regionInfo) {}, func(context.Context, rangeTask) {}) + key := newRegionRecoveryKey(span.subID, span.span) + failureHandler.recoveries[key] = ®ionRecoveryState{} + + region := newRegionInfo(tikv.NewRegionVerID(1, 1, 1), span.span, nil, span, false) + region.lockedRangeState = ®ionlock.LockedRangeState{} + region.rpcCtx = &tikv.RPCContext{Addr: "store-1"} + state := newRegionFeedState(region, uint64(span.subID), ®ionRequestWorker{tracker: newRegionTracker()}, nil, func(state *regionFeedState) { + failureHandler.resetRegionRecovery(state.region) + }) + + handler := ®ionEventHandler{ + eventSink: ®ionEventSink{memoryQuota: newMemoryQuotaController(0, 0)}, + failureHandler: failureHandler, + } + handler.Handle(span, regionEvent{ + states: []*regionFeedState{state}, + entries: &cdcpb.Event_Entries_{ + Entries: &cdcpb.Event_Entries{ + Entries: []*cdcpb.Event_Row{{Type: cdcpb.Event_INITIALIZED}}, + }, + }, + }) + + failureHandler.recoveryMu.Lock() + _, ok := failureHandler.recoveries[key] + failureHandler.recoveryMu.Unlock() + require.False(t, ok) +} + func TestSpanInitializedAfterFullRangeCoverage(t *testing.T) { const startTs = 100 span := &subscribedSpan{ @@ -448,7 +485,7 @@ func TestSpanInitializedAfterFullRangeCoverage(t *testing.T) { }, subscribedSpan: span, lockedRangeState: ®ionlock.LockedRangeState{}, - }, uint64(span.subID), ®ionRequestWorker{}, nil) + }, uint64(span.subID), ®ionRequestWorker{}, nil, nil) secondState := newRegionFeedState(regionInfo{ verID: tikv.NewRegionVerID(2, 1, 1), span: heartbeatpb.TableSpan{ @@ -457,7 +494,7 @@ func TestSpanInitializedAfterFullRangeCoverage(t *testing.T) { }, subscribedSpan: span, lockedRangeState: ®ionlock.LockedRangeState{}, - }, uint64(span.subID), ®ionRequestWorker{}, nil) + }, uint64(span.subID), ®ionRequestWorker{}, nil, nil) span.markRegionInitialized(firstState) require.False(t, span.initialized.Load()) diff --git a/logservice/logpuller/region_failure_handler.go b/logservice/logpuller/region_failure_handler.go index 3622d927d9..64568df09c 100644 --- a/logservice/logpuller/region_failure_handler.go +++ b/logservice/logpuller/region_failure_handler.go @@ -15,11 +15,13 @@ package logpuller import ( "context" + "math/rand/v2" "sync" "time" "github.com/pingcap/kvproto/pkg/cdcpb" "github.com/pingcap/log" + "github.com/pingcap/ticdc/heartbeatpb" "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/metrics" "github.com/tikv/client-go/v2/tikv" @@ -43,12 +45,60 @@ var ( type regionFailureHandler struct { cache *errCache regionCache *tikv.RegionCache + recoveryMu sync.Mutex + recoveries map[regionRecoveryKey]*regionRecoveryState onTableDrained func(*subscribedSpan) scheduleRegionRequest func(context.Context, regionInfo) scheduleRangeRequest func(context.Context, rangeTask) } +const ( + regionRecoveryBaseDelay = 50 * time.Millisecond + regionRecoveryMaxDelay = 2 * time.Second + regionRecoveryStateTTL = 5 * time.Minute +) + +// regionRecoveryKey keeps backoff state across region ID and epoch changes for +// the same logical range. +type regionRecoveryKey struct { + subscriptionID SubscriptionID + startKey string + endKey string +} + +type regionRecoveryState struct { + attempt uint32 + expiresAt time.Time +} + +func newRegionRecoveryKey( + subscriptionID SubscriptionID, + span heartbeatpb.TableSpan, +) regionRecoveryKey { + return regionRecoveryKey{ + subscriptionID: subscriptionID, + startKey: string(span.StartKey), + endKey: string(span.EndKey), + } +} + +func regionRecoveryDelay(attempt uint32) time.Duration { + if attempt == 0 { + attempt = 1 + } + exponent := attempt - 1 + if exponent > 16 { + exponent = 16 + } + delay := regionRecoveryBaseDelay << exponent + if delay > regionRecoveryMaxDelay { + delay = regionRecoveryMaxDelay + } + half := delay / 2 + return half + time.Duration(rand.Int64N(int64(delay-half)+1)) +} + func newRegionFailureHandler( regionCache *tikv.RegionCache, onTableDrained func(*subscribedSpan), @@ -58,12 +108,85 @@ func newRegionFailureHandler( return ®ionFailureHandler{ cache: newErrCache(), regionCache: regionCache, + recoveries: make(map[regionRecoveryKey]*regionRecoveryState), onTableDrained: onTableDrained, scheduleRegionRequest: scheduleRegionRequest, scheduleRangeRequest: scheduleRangeRequest, } } +func (r *regionFailureHandler) scheduleRecovery( + ctx context.Context, + subscribedSpan *subscribedSpan, + span heartbeatpb.TableSpan, + minDelay time.Duration, + retry func(), +) { + if subscribedSpan == nil || subscribedSpan.stopped.Load() { + return + } + key := newRegionRecoveryKey(subscribedSpan.subID, span) + + r.recoveryMu.Lock() + state := r.recoveries[key] + if state == nil { + state = ®ionRecoveryState{} + r.recoveries[key] = state + } + if state.attempt < 32 { + state.attempt++ + } + delay := regionRecoveryDelay(state.attempt) + if minDelay > delay { + delay = minDelay + } + state.expiresAt = time.Now().Add(delay + regionRecoveryStateTTL) + r.recoveryMu.Unlock() + + time.AfterFunc(delay, func() { + r.recoveryMu.Lock() + if r.recoveries[key] != state { + r.recoveryMu.Unlock() + return + } + // Keep the attempt until the retry succeeds or the state expires. + state.expiresAt = time.Now().Add(regionRecoveryStateTTL) + r.recoveryMu.Unlock() + + if ctx.Err() != nil || subscribedSpan.stopped.Load() { + r.resetRecovery(key) + return + } + retry() + }) +} + +func (r *regionFailureHandler) expireRecoveries(now time.Time) { + r.recoveryMu.Lock() + defer r.recoveryMu.Unlock() + for key, state := range r.recoveries { + if !state.expiresAt.After(now) { + delete(r.recoveries, key) + } + } +} + +func (r *regionFailureHandler) resetRecovery(key regionRecoveryKey) { + r.recoveryMu.Lock() + defer r.recoveryMu.Unlock() + delete(r.recoveries, key) +} + +func (r *regionFailureHandler) resetRegionRecovery(region regionInfo) { + r.resetRecovery(newRegionRecoveryKey(region.subscribedSpan.subID, region.span)) +} + +func (r *regionFailureHandler) cancelRecoveries() { + r.recoveryMu.Lock() + defer r.recoveryMu.Unlock() + clear(r.recoveries) +} + // Report admits a region failure into the recovery pipeline. It releases the // corresponding range lock before enqueueing the failure so new range tasks are // not blocked by stale region ownership. @@ -80,6 +203,7 @@ func (r *regionFailureHandler) Report(errInfo regionErrorInfo) { func (r *regionFailureHandler) Run(ctx context.Context) error { log.Info("region failure handler starts") defer log.Info("region failure handler exits") + defer r.cancelRecoveries() handleCachedErrors := func() error { for { @@ -102,13 +226,17 @@ func (r *regionFailureHandler) Run(ctx context.Context) error { // r.cache.ready() should handle failures promptly in normal flow. The ticker is only a // fallback scan and is not expected to be needed in practice. - ticker := time.NewTicker(200 * time.Millisecond) - defer ticker.Stop() + fallbackTicker := time.NewTicker(200 * time.Millisecond) + defer fallbackTicker.Stop() + cleanupTicker := time.NewTicker(regionRecoveryStateTTL) + defer cleanupTicker.Stop() for { select { case <-ctx.Done(): return ctx.Err() - case <-ticker.C: + case now := <-cleanupTicker.C: + r.expireRecoveries(now) + case <-fallbackTicker.C: if err := handleCachedErrors(); err != nil { return err } @@ -122,7 +250,18 @@ func (r *regionFailureHandler) Run(ctx context.Context) error { func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionErrorInfo) error { err := errors.Cause(errInfo.err) - rescheduleRange := func() { + retryRegion := func(minDelay time.Duration) { + r.scheduleRecovery( + ctx, + errInfo.subscribedSpan, + errInfo.span, + minDelay, + func() { + r.scheduleRegionRequest(ctx, errInfo.regionInfo) + }, + ) + } + retryRange := func() { priority := normalizeScanPriority(errInfo.scanPriority) if priority == cdcpb.ScanPriority_SCAN_PRIORITY_LOW { priority = errInfo.subscribedSpan.priorityPolicy.resolve( @@ -131,12 +270,21 @@ func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionEr errInfo.subscribedSpan.priorityPolicy.pdClock.CurrentTime(), ) } - r.scheduleRangeRequest(ctx, rangeTask{ + task := rangeTask{ span: errInfo.span, subscribedSpan: errInfo.subscribedSpan, filterLoop: errInfo.filterLoop, priority: priority, - }) + } + r.scheduleRecovery( + ctx, + task.subscribedSpan, + task.span, + 0, + func() { + r.scheduleRangeRequest(ctx, task) + }, + ) } //nolint:errorlint // converting large type switch to errors.As is a significant refactor @@ -153,28 +301,34 @@ func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionEr innerErr := eerr.err if notLeader := innerErr.GetNotLeader(); notLeader != nil { metricFeedNotLeaderCounter.Inc() - r.regionCache.UpdateLeader(errInfo.verID, notLeader.GetLeader(), errInfo.rpcCtx.AccessIdx) - r.scheduleRegionRequest(ctx, errInfo.regionInfo) + leader := notLeader.GetLeader() + if leader == nil || leader.GetId() == 0 || leader.GetStoreId() == 0 || errInfo.rpcCtx == nil { + r.regionCache.InvalidateCachedRegion(errInfo.verID) + retryRange() + return nil + } + r.regionCache.UpdateLeader(errInfo.verID, leader, errInfo.rpcCtx.AccessIdx) + retryRegion(0) return nil } if innerErr.GetEpochNotMatch() != nil { metricFeedEpochNotMatchCounter.Inc() - rescheduleRange() + retryRange() return nil } if innerErr.GetRegionNotFound() != nil { metricFeedRegionNotFoundCounter.Inc() - rescheduleRange() + retryRange() return nil } if innerErr.GetCongested() != nil { metricKvCongestedCounter.Inc() - r.scheduleRegionRequest(ctx, errInfo.regionInfo) + retryRegion(0) return nil } - if innerErr.GetServerIsBusy() != nil { + if busy := innerErr.GetServerIsBusy(); busy != nil { metricKvIsBusyCounter.Inc() - r.scheduleRegionRequest(ctx, errInfo.regionInfo) + retryRegion(time.Duration(busy.GetBackoffMs()) * time.Millisecond) return nil } if duplicated := innerErr.GetDuplicateRequest(); duplicated != nil { @@ -193,27 +347,30 @@ func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionEr zap.Uint64("subscriptionID", uint64(errInfo.subscribedSpan.subID)), zap.Stringer("error", innerErr)) metricFeedUnknownErrorCounter.Inc() - r.scheduleRegionRequest(ctx, errInfo.regionInfo) + retryRegion(0) return nil case *rpcCtxUnavailableErr: metricFeedRPCCtxUnavailable.Inc() - rescheduleRange() + retryRange() return nil case *getStoreErr: metricGetStoreErr.Inc() bo := tikv.NewBackoffer(ctx, tikvRequestMaxBackoff) // cannot get the store the region belongs to, so we need to reload the region. r.regionCache.OnSendFail(bo, errInfo.rpcCtx, true, err) - rescheduleRange() + retryRange() return nil case *storeStreamErr: metricStoreSendRequestErr.Inc() bo := tikv.NewBackoffer(ctx, tikvRequestMaxBackoff) r.regionCache.OnSendFail(bo, errInfo.rpcCtx, regionScheduleReload, err) - r.scheduleRegionRequest(ctx, errInfo.regionInfo) + retryRegion(0) return nil case *requestCancelledErr: // the corresponding subscription has been unsubscribed, just ignore. + if errInfo.subscribedSpan != nil { + r.resetRegionRecovery(errInfo.regionInfo) + } return nil default: // TODO(qupeng): for some errors it's better to just deregister the region from TiKVs. diff --git a/logservice/logpuller/region_failure_handler_test.go b/logservice/logpuller/region_failure_handler_test.go index f47fca3509..0ce52a44a5 100644 --- a/logservice/logpuller/region_failure_handler_test.go +++ b/logservice/logpuller/region_failure_handler_test.go @@ -19,7 +19,11 @@ import ( "time" "github.com/pingcap/errors" + "github.com/pingcap/kvproto/pkg/cdcpb" + "github.com/pingcap/kvproto/pkg/errorpb" "github.com/pingcap/ticdc/heartbeatpb" + "github.com/pingcap/ticdc/pkg/pdutil" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/tikv/client-go/v2/tikv" ) @@ -142,3 +146,118 @@ func TestRegionFailureHandlerRunDrainsErrCacheWithoutDispatcher(t *testing.T) { t.Fatal("failure handler did not exit after context cancellation") } } + +func TestRegionFailureHandlerSchedulesNotLeaderRangeRetry(t *testing.T) { + pdClient := newFailureRecoveryTestPDClient(t) + defer pdClient.Close() + + regionCache := tikv.NewRegionCache(pdClient) + defer regionCache.Close() + + region := createFailureRecoveryTestRegion(t, SubscriptionID(1), 1) + region.subscribedSpan.priorityPolicy = newScanPriorityPolicy(pdutil.NewClock4Test(), 30*time.Minute) + + rangeRetryCh := make(chan rangeTask, 2) + handler := newRegionFailureHandler( + regionCache, + func(*subscribedSpan) {}, + func(context.Context, regionInfo) { + t.Fatal("unexpected region retry") + }, + func(_ context.Context, task rangeTask) { + rangeRetryCh <- task + }, + ) + errInfo := newRegionErrorInfo(region, &eventError{ + err: &cdcpb.Error{NotLeader: &errorpb.NotLeader{}}, + }) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + require.NoError(t, handler.handleError(ctx, errInfo)) + + select { + case task := <-rangeRetryCh: + require.Equal(t, region.span, task.span) + require.Same(t, region.subscribedSpan, task.subscribedSpan) + case <-time.After(time.Second): + t.Fatal("not leader retry was not scheduled") + } +} + +func TestRegionRecoveryBackoffFollowsRangeAcrossRegionChanges(t *testing.T) { + regionRetryCh := make(chan regionInfo, 2) + handler := newRegionFailureHandler( + nil, + func(*subscribedSpan) {}, + func(_ context.Context, region regionInfo) { + regionRetryCh <- region + }, + func(context.Context, rangeTask) {}, + ) + t.Cleanup(handler.cancelRecoveries) + + region := createFailureRecoveryTestRegion(t, SubscriptionID(1), 1) + errInfo := newRegionErrorInfo(region, &eventError{ + err: &cdcpb.Error{Congested: &cdcpb.Congested{}}, + }) + require.NoError(t, handler.handleError(context.Background(), errInfo)) + select { + case retried := <-regionRetryCh: + require.Equal(t, uint64(1), retried.verID.GetID()) + case <-time.After(time.Second): + t.Fatal("first region recovery was not scheduled") + } + + region.verID = tikv.NewRegionVerID(2, 1, 1) + errInfo = newRegionErrorInfo(region, &eventError{ + err: &cdcpb.Error{Congested: &cdcpb.Congested{}}, + }) + require.NoError(t, handler.handleError(context.Background(), errInfo)) + select { + case retried := <-regionRetryCh: + require.Equal(t, uint64(2), retried.verID.GetID()) + case <-time.After(time.Second): + t.Fatal("second region recovery was not scheduled") + } + + key := newRegionRecoveryKey(region.subscribedSpan.subID, region.span) + handler.recoveryMu.Lock() + attempt := handler.recoveries[key].attempt + handler.recoveryMu.Unlock() + require.Equal(t, uint32(2), attempt) +} + +func TestRegionFailureHandlerRequestCancelledResetsRecoveryState(t *testing.T) { + handler := newRegionFailureHandler(nil, func(*subscribedSpan) {}, func(context.Context, regionInfo) {}, func(context.Context, rangeTask) {}) + region := createFailureRecoveryTestRegion(t, SubscriptionID(1), 1) + key := newRegionRecoveryKey(region.subscribedSpan.subID, region.span) + handler.recoveries[key] = ®ionRecoveryState{} + + err := handler.handleError(context.Background(), newRegionErrorInfo(region, &requestCancelledErr{})) + require.NoError(t, err) + + handler.recoveryMu.Lock() + _, ok := handler.recoveries[key] + handler.recoveryMu.Unlock() + assert.False(t, ok) +} + +func TestRegionFailureHandlerExpiresRecoveryStates(t *testing.T) { + handler := newRegionFailureHandler(nil, func(*subscribedSpan) {}, func(context.Context, regionInfo) {}, func(context.Context, rangeTask) {}) + now := time.Now() + expiredKey := newRegionRecoveryKey(1, heartbeatpb.TableSpan{StartKey: []byte("a"), EndKey: []byte("b")}) + activeKey := newRegionRecoveryKey(1, heartbeatpb.TableSpan{StartKey: []byte("b"), EndKey: []byte("c")}) + handler.recoveries[expiredKey] = ®ionRecoveryState{expiresAt: now.Add(-time.Second)} + handler.recoveries[activeKey] = ®ionRecoveryState{expiresAt: now.Add(time.Second)} + + handler.expireRecoveries(now) + + handler.recoveryMu.Lock() + _, expiredExists := handler.recoveries[expiredKey] + _, activeExists := handler.recoveries[activeKey] + handler.recoveryMu.Unlock() + require.False(t, expiredExists) + require.True(t, activeExists) +} diff --git a/logservice/logpuller/region_request_worker.go b/logservice/logpuller/region_request_worker.go index 3f954168d0..614b67d7f0 100644 --- a/logservice/logpuller/region_request_worker.go +++ b/logservice/logpuller/region_request_worker.go @@ -454,7 +454,9 @@ func (s *regionRequestWorker) sendRegionRequest(conn *ConnAndClient, req *region // Publish the state before Send so a fast response observes its owner and // admission lease. - state := newRegionFeedState(region, uint64(subID), s, req) + state := newRegionFeedState(region, uint64(subID), s, req, func(state *regionFeedState) { + s.failureHandler.resetRegionRecovery(state.region) + }) if !s.tracker.Add(subID, region.verID.GetID(), state) { // RangeLock normally prevents duplicate active regions. Keep the existing // owner, including its range-lock ownership, if that invariant is ever diff --git a/logservice/logpuller/region_request_worker_test.go b/logservice/logpuller/region_request_worker_test.go index 59da548746..4f429925ee 100644 --- a/logservice/logpuller/region_request_worker_test.go +++ b/logservice/logpuller/region_request_worker_test.go @@ -217,7 +217,7 @@ func TestRegionRequestWorkerIgnoresDuplicateActiveRegion(t *testing.T) { region := prepareRegionForSendTest(createTestRegionInfo(1, 1)) req1 := admitRegionRequest(t, admission, region) - state1 := newRegionFeedState(region, uint64(region.subscribedSpan.subID), worker, req1) + state1 := newRegionFeedState(region, uint64(region.subscribedSpan.subID), worker, req1, nil) require.True(t, worker.tracker.Add(region.subscribedSpan.subID, region.verID.GetID(), state1)) req2 := admitRegionRequest(t, admission, region) @@ -456,7 +456,7 @@ func TestStoppedStateRemovesSentRequest(t *testing.T) { region := prepareRegionForSendTest(createTestRegionInfo(1, 1)) req := admitRegionRequest(t, admission, region) - state := newRegionFeedState(req.regionInfo, uint64(req.regionInfo.subscribedSpan.subID), worker, req) + state := newRegionFeedState(req.regionInfo, uint64(req.regionInfo.subscribedSpan.subID), worker, req, nil) require.True(t, worker.tracker.Add(req.regionInfo.subscribedSpan.subID, req.regionInfo.verID.GetID(), state)) state.markStopped(errors.New("send request to store error")) worker.tracker.RemoveIf(req.regionInfo.subscribedSpan.subID, req.regionInfo.verID.GetID(), state) @@ -482,7 +482,7 @@ func TestRunStreamFailurePushesTrackedRegionToEventSink(t *testing.T) { sentRegion := createFailureRecoveryTestRegion(t, 1, 1) sentReq := admitRegionRequest(t, worker.admission, sentRegion) - sentState := newRegionFeedState(sentRegion, uint64(sentRegion.subscribedSpan.subID), worker, sentReq) + sentState := newRegionFeedState(sentRegion, uint64(sentRegion.subscribedSpan.subID), worker, sentReq, nil) require.True(t, worker.tracker.Add(sentRegion.subscribedSpan.subID, sentRegion.verID.GetID(), sentState)) firstRegion := createFailureRecoveryTestRegion(t, 2, 2) diff --git a/logservice/logpuller/region_state.go b/logservice/logpuller/region_state.go index 0f2f84ce3e..422c3bb39c 100644 --- a/logservice/logpuller/region_state.go +++ b/logservice/logpuller/region_state.go @@ -90,6 +90,9 @@ type regionFeedState struct { region regionInfo requestID uint64 // It is also the subscription ID matcher *matcher + // onInitialized runs once when the region finishes its first successful + // initialization for this request lifecycle. + onInitialized func(*regionFeedState) // Transform: normal -> stopped -> removed. // normal: the region is in replicating. @@ -113,12 +116,14 @@ func newRegionFeedState( requestID uint64, worker *regionRequestWorker, request *regionReq, + onInitialized func(*regionFeedState), ) *regionFeedState { state := ®ionFeedState{ - region: region, - requestID: requestID, - matcher: newMatcher(), - worker: worker, + region: region, + requestID: requestID, + matcher: newMatcher(), + onInitialized: onInitialized, + worker: worker, } state.regionReq.Store(request) return state @@ -167,8 +172,13 @@ func (s *regionFeedState) isInitialized() bool { } func (s *regionFeedState) setInitialized() { - s.region.lockedRangeState.Initialized.Store(true) + if !s.region.lockedRangeState.Initialized.CompareAndSwap(false, true) { + return + } s.finishScan() + if s.onInitialized != nil { + s.onInitialized(s) + } } func (s *regionFeedState) finishScan() { diff --git a/logservice/logpuller/subscription_client_test.go b/logservice/logpuller/subscription_client_test.go index c7e4591f92..fa42a096c1 100644 --- a/logservice/logpuller/subscription_client_test.go +++ b/logservice/logpuller/subscription_client_test.go @@ -57,6 +57,7 @@ func TestGenerateResolveLockTask(t *testing.T) { client := &subscriptionClient{ resolveLockTaskCh: make(chan resolveLockTask, 10), resolveLockRateLimiter: newResolveLockRateLimiter(), + memoryQuota: newMemoryQuotaController(0, 0), } client.ctx, client.cancel = context.WithCancel(context.Background()) rawSpan := heartbeatpb.TableSpan{ @@ -111,7 +112,7 @@ func TestGenerateResolveLockTask(t *testing.T) { // Lock another range, no task will be triggered before initialized. res = span.rangeLock.LockRange(context.Background(), []byte{'c'}, []byte{'d'}, 2, 100) require.Equal(t, regionlock.LockRangeStatusSuccess, res.Status) - state := newRegionFeedState(regionInfo{lockedRangeState: res.LockedRangeState, subscribedSpan: span}, 1, worker, nil) + state := newRegionFeedState(regionInfo{lockedRangeState: res.LockedRangeState, subscribedSpan: span}, 1, worker, nil, nil) span.resolveStaleLocks(200) select { case <-client.resolveLockTaskCh: @@ -305,7 +306,9 @@ func TestResolveLockTaskDroppedWhenChannelFull(t *testing.T) { func TestStopTaskUsesSubscribedSpanFilterLoop(t *testing.T) { client := &subscriptionClient{ - resolveLockTaskCh: make(chan resolveLockTask, 1), + resolveLockTaskCh: make(chan resolveLockTask, 1), + resolveLockRateLimiter: newResolveLockRateLimiter(), + memoryQuota: newMemoryQuotaController(0, 0), } client.ctx, client.cancel = context.WithCancel(context.Background()) defer client.cancel() @@ -457,7 +460,6 @@ func TestRegionEventSinkPushUnblocksOnClientClose(t *testing.T) { }, }, } - done := make(chan struct{}) go func() { sink.Push(SubscriptionID(1), event)