Skip to content
Draft
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
2 changes: 1 addition & 1 deletion logservice/eventstore/pebble.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
12 changes: 5 additions & 7 deletions logservice/logpuller/region_admission_controller_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
49 changes: 43 additions & 6 deletions logservice/logpuller/region_event_handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ func TestHandleEventEntryEventOutOfOrder(t *testing.T) {
false,
)
region.lockedRangeState = &regionlock.LockedRangeState{}
state := newRegionFeedState(region, 1, worker, nil)
state := newRegionFeedState(region, 1, worker, nil, nil)

// Receive prewrite2 with empty value.
{
Expand Down Expand Up @@ -225,7 +225,7 @@ func TestHandleResolvedTs(t *testing.T) {
worker := &regionRequestWorker{
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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -378,6 +378,7 @@ func TestHandleResolvedTsThrottled(t *testing.T) {
1,
worker,
nil,
nil,
)

require.Equal(t, uint64(200), handleResolvedTs(span, state, 300))
Expand Down Expand Up @@ -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] = &regionRecoveryState{}

region := newRegionInfo(tikv.NewRegionVerID(1, 1, 1), span.span, nil, span, false)
region.lockedRangeState = &regionlock.LockedRangeState{}
region.rpcCtx = &tikv.RPCContext{Addr: "store-1"}
state := newRegionFeedState(region, uint64(span.subID), &regionRequestWorker{tracker: newRegionTracker()}, nil, func(state *regionFeedState) {
failureHandler.resetRegionRecovery(state.region)
})

handler := &regionEventHandler{
eventSink: &regionEventSink{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{
Expand All @@ -448,7 +485,7 @@ func TestSpanInitializedAfterFullRangeCoverage(t *testing.T) {
},
subscribedSpan: span,
lockedRangeState: &regionlock.LockedRangeState{},
}, uint64(span.subID), &regionRequestWorker{}, nil)
}, uint64(span.subID), &regionRequestWorker{}, nil, nil)
secondState := newRegionFeedState(regionInfo{
verID: tikv.NewRegionVerID(2, 1, 1),
span: heartbeatpb.TableSpan{
Expand All @@ -457,7 +494,7 @@ func TestSpanInitializedAfterFullRangeCoverage(t *testing.T) {
},
subscribedSpan: span,
lockedRangeState: &regionlock.LockedRangeState{},
}, uint64(span.subID), &regionRequestWorker{}, nil)
}, uint64(span.subID), &regionRequestWorker{}, nil, nil)

span.markRegionInitialized(firstState)
require.False(t, span.initialized.Load())
Expand Down
Loading