From 593f2101d2328c120ab8b8dafe09fc9e5ea246f5 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Mon, 27 Jul 2026 08:51:44 +0000 Subject: [PATCH] init Signed-off-by: wk989898 --- coordinator/changefeed/changefeed.go | 15 ++++++++- coordinator/changefeed/changefeed_test.go | 39 +++++++++++++++++++++++ 2 files changed, 53 insertions(+), 1 deletion(-) diff --git a/coordinator/changefeed/changefeed.go b/coordinator/changefeed/changefeed.go index f75c7f9754..f840835d90 100644 --- a/coordinator/changefeed/changefeed.go +++ b/coordinator/changefeed/changefeed.go @@ -166,12 +166,25 @@ func (c *Changefeed) ShouldRun() bool { // It returns false if the status is not changed // It returns the new state and error if the status is changed func (c *Changefeed) UpdateStatus(newStatus *heartbeatpb.MaintainerStatus) (bool, config.FeedState, *heartbeatpb.RunningError) { + if newStatus == nil { + return false, config.StateNormal, nil + } + old := c.status.Load() failpoint.Inject("CoordinatorDontUpdateChangefeedCheckpoint", func() { newStatus.CheckpointTs = old.CheckpointTs }) - if newStatus != nil && newStatus.CheckpointTs >= old.CheckpointTs { + if newStatus.CheckpointTs < old.CheckpointTs { + if len(newStatus.Err) == 0 { + return false, config.StateNormal, nil + } + statusWithMonotonicCheckpoint := *newStatus + statusWithMonotonicCheckpoint.CheckpointTs = old.CheckpointTs + newStatus = &statusWithMonotonicCheckpoint + } + + if newStatus.CheckpointTs >= old.CheckpointTs { c.status.Store(newStatus) changed, state, err := c.backoff.checkFailedStatus(newStatus) diff --git a/coordinator/changefeed/changefeed_test.go b/coordinator/changefeed/changefeed_test.go index f9652ed325..3c106029b5 100644 --- a/coordinator/changefeed/changefeed_test.go +++ b/coordinator/changefeed/changefeed_test.go @@ -110,6 +110,45 @@ func TestChangefeed_UpdateStatus(t *testing.T) { require.Equal(t, newStatus, cf.GetStatus()) } +func TestChangefeed_UpdateStatusProcessesErrorsWhenCheckpointRegresses(t *testing.T) { + cfID := common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName) + info := &config.ChangeFeedInfo{ + SinkURI: "kafka://127.0.0.1:9092", + State: config.StateNormal, + Config: config.GetDefaultReplicaConfig(), + } + cf := NewChangefeed(cfID, info, 200, true) + + regressedStatus := &heartbeatpb.MaintainerStatus{CheckpointTs: 150} + updated, state, err := cf.UpdateStatus(regressedStatus) + require.False(t, updated) + require.Equal(t, config.StateNormal, state) + require.Nil(t, err) + require.Equal(t, uint64(200), cf.GetStatus().CheckpointTs) + + retryableErr := &heartbeatpb.RunningError{ + Node: "node-1", + Code: "CDC:ErrChangefeedRetryable", + Message: "retryable error", + } + regressedStatus = &heartbeatpb.MaintainerStatus{ + CheckpointTs: 150, + Err: []*heartbeatpb.RunningError{retryableErr}, + } + updated, state, err = cf.UpdateStatus(regressedStatus) + require.True(t, updated) + require.Equal(t, config.StateWarning, state) + require.Same(t, retryableErr, err) + require.Equal(t, uint64(150), regressedStatus.CheckpointTs) + status := cf.GetStatus() + require.Equal(t, uint64(200), status.CheckpointTs) + require.Len(t, status.Err, 1) + require.Same(t, retryableErr, status.Err[0]) + require.False(t, cf.ShouldRun()) + require.True(t, cf.backoff.retrying.Load()) + require.True(t, cf.backoff.isRestarting.Load()) +} + func TestChangefeed_UpdateStatusFastFailWhenBootstrapDoneChanges(t *testing.T) { cfID := common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName) info := &config.ChangeFeedInfo{