Skip to content
Open
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
11 changes: 6 additions & 5 deletions pkg/eventservice/event_broker.go
Original file line number Diff line number Diff line change
Expand Up @@ -522,21 +522,22 @@ func (c *eventBroker) getScanTaskRequestResult(task scanTask) scanTaskRequestRes
}
}

hasRowResume := len(request.Cursor.Position) != 0
// A published row cursor at C came from an earlier scan whose DDL and received
hasResumeCursor := request.Cursor.TxnStartTs != 0 || len(request.Cursor.Position) != 0
// A published cursor at C came from an earlier scan whose DDL and received
// resolved-ts bounds had already reached C. Since those bounds do not regress,
// only the adaptive scan window can move CommitTsEnd behind C. For example, if
// C=100 and the window caps the end at 80, restore the effective range to
// [100, 100] so scanning resumes after Position inside that transaction.
if hasRowResume && dataRange.CommitTsEnd < dataRange.CommitTsStart {
// [100, 100] so Position can resume rows inside a transaction, or TxnStartTs
// can resume later transactions sharing commit-ts C.
if hasResumeCursor && dataRange.CommitTsEnd < dataRange.CommitTsStart {
dataRange.CommitTsEnd = dataRange.CommitTsStart
}

if dataRange.CommitTsEnd <= dataRange.CommitTsStart {
// A cursor makes [C, C] meaningful: Position resumes rows inside a
// transaction, while TxnStartTs resumes later transactions at the same C.
canResumeAtStart := dataRange.CommitTsEnd == dataRange.CommitTsStart &&
(hasRowResume || request.Cursor.TxnStartTs != 0)
hasResumeCursor
if canResumeAtStart || task.hasPendingLargeTxnState() {
result := scanTaskRequestResult{needScan: true, request: request}
if task.changefeedStat.lowLatencyMode {
Expand Down
10 changes: 7 additions & 3 deletions pkg/eventservice/event_broker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -760,7 +760,7 @@ func TestScanRangeCappedByScanWindow(t *testing.T) {
require.Equal(t, oracle.GoTimeToTS(baseTime.Add(defaultScanInterval)), result.request.Range.CommitTsEnd)
}

func TestGetScanTaskDataRangeEmptyAfterCappingDoesNotResetScanRange(t *testing.T) {
func TestGetScanTaskRequestKeepsTxnCursorInsideShrunkWindow(t *testing.T) {
broker, _, _, _ := newEventBrokerForTest()
// Close the broker, so we can catch all message in the test.
broker.close()
Expand All @@ -785,8 +785,12 @@ func TestGetScanTaskDataRangeEmptyAfterCappingDoesNotResetScanRange(t *testing.T
changefeedStatus.minSentTs.Store(baseTs)
changefeedStatus.scanInterval.Store(int64(defaultScanInterval))

needScan, _ := broker.getScanTaskRequest(disp)
require.False(t, needScan)
needScan, request := broker.getScanTaskRequest(disp)
require.True(t, needScan)
require.Equal(t, commitStart, request.Range.CommitTsStart)
require.Equal(t, commitStart, request.Range.CommitTsEnd)
require.Equal(t, lastStartTs, request.Cursor.TxnStartTs)
require.Empty(t, request.Cursor.Position)
require.Equal(t, commitStart, disp.loadScanProgress().txnCommitTs)
require.Equal(t, lastStartTs, disp.loadScanProgress().txnStartTs)
}
Expand Down
Loading