Skip to content
Merged
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
29 changes: 24 additions & 5 deletions coordinator/operator/operator_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -285,7 +285,8 @@ func (oc *Controller) StopRemoteMaintainerWithMaintainerEpoch(
if changefeed != nil {
keyspaceID = changefeed.GetKeyspaceID()
}
return oc.pushStopChangefeedOperator(keyspaceID, cfID, nodeID, removed, maintainerEpoch)
return oc.pushStopChangefeedOperator(
keyspaceID, cfID, nodeID, removed, maintainerEpoch, stopChangefeedKindStaleOwner)
}

// stopChangefeed creates a stop operator using the owner epoch that must be
Expand All @@ -299,6 +300,15 @@ func (oc *Controller) stopChangefeed(
oc.mu.Lock()
defer oc.mu.Unlock()

if !removed {
if old, ok := oc.operators[cfID]; ok {
if oldStop, ok := old.OP.(*StopChangefeedOperator); ok &&
oldStop.kind == stopChangefeedKindCurrentPlacement {
return oldStop
}
}
}

changefeed := oc.changefeedDB.GetByID(cfID)
keyspaceID := common.DefaultKeyspaceID
if changefeed != nil {
Expand Down Expand Up @@ -328,7 +338,8 @@ func (oc *Controller) stopChangefeed(
scheduledNode = oc.selfNode.ID
}

return oc.pushStopChangefeedOperator(keyspaceID, cfID, scheduledNode, removed, maintainerEpoch)
return oc.pushStopChangefeedOperator(
keyspaceID, cfID, scheduledNode, removed, maintainerEpoch, stopChangefeedKindCurrentPlacement)
}

// moveOriginStopTargetLocked returns the origin maintainer stop target for an
Expand All @@ -346,16 +357,16 @@ func (oc *Controller) moveOriginStopTargetLocked(cfID common.ChangeFeedID) (node
}

// pushStopChangefeedOperator pushes a stop changefeed operator to the controller.
// it checks if the operator already exists, if exists, it will replace the old one.
// if the old operator is the removing operator, it will skip this operator.
// Keep an existing placement stop for repeated non-removing requests. A stale-owner
// cleanup does not represent the current placement and may be replaced.
func (oc *Controller) pushStopChangefeedOperator(
keyspaceID uint32,
cfID common.ChangeFeedID,
nodeID node.ID,
remove bool,
maintainerEpoch uint64,
kind stopChangefeedKind,
) operator.Operator[common.ChangeFeedID, *heartbeatpb.MaintainerStatus] {
op := NewStopChangefeedOperator(keyspaceID, cfID, nodeID, oc.selfNode.ID, oc.backend, remove, maintainerEpoch)
if old, ok := oc.operators[cfID]; ok {
oldStop, ok := old.OP.(*StopChangefeedOperator)
if ok {
Expand All @@ -365,6 +376,13 @@ func (oc *Controller) pushStopChangefeedOperator(
zap.String("changefeed", cfID.Name()))
return oldStop
}
if !remove {
if oldStop.kind == stopChangefeedKindCurrentPlacement ||
(kind == stopChangefeedKindStaleOwner &&
oldStop.nodeID == nodeID && oldStop.maintainerEpoch == maintainerEpoch) {
return oldStop
}
}
}
log.Info("changefeed is stopped, replace the old one",
zap.String("role", oc.role),
Expand All @@ -375,6 +393,7 @@ func (oc *Controller) pushStopChangefeedOperator(
old.IsRemoved.Store(true)
delete(oc.operators, old.OP.ID())
}
op := NewStopChangefeedOperator(keyspaceID, cfID, nodeID, oc.selfNode.ID, oc.backend, remove, maintainerEpoch, kind)
oc.pushOperator(op)
return op
}
Expand Down
52 changes: 49 additions & 3 deletions coordinator/operator/operator_controller_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,9 @@ func TestController_StopChangefeedWithMaintainerEpoch(t *testing.T) {
changefeedDB := changefeed.NewChangefeedDB(1216)
ctrl := gomock.NewController(t)
backend := mock_changefeed.NewMockBackend(ctrl)
oc, self, _ := newOperatorControllerForTest(t, changefeedDB, backend, nil)
oc, _, nodeManager := newOperatorControllerForTest(t, changefeedDB, backend, nil)
owner := node.NewInfo("localhost:8301", "")
nodeManager.GetAliveNodes()[owner.ID] = owner
cfID := common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName)
cf := changefeed.NewChangefeed(cfID, &config.ChangeFeedInfo{
ChangefeedID: cfID,
Expand All @@ -101,11 +103,24 @@ func TestController_StopChangefeedWithMaintainerEpoch(t *testing.T) {
Epoch: 20,
},
1, true)
changefeedDB.AddReplicatingMaintainer(cf, self.ID)
changefeedDB.AddReplicatingMaintainer(cf, owner.ID)

op := oc.StopChangefeedWithMaintainerEpoch(context.Background(), cfID, false, 10)
req := op.Schedule().Message[0].(*heartbeatpb.RemoveMaintainerRequest)
reqMsg := op.Schedule()
require.Equal(t, owner.ID, reqMsg.To)
req := reqMsg.Message[0].(*heartbeatpb.RemoveMaintainerRequest)
require.Equal(t, uint64(10), req.MaintainerEpoch)

// A repeated warning arrives after the first stop has cleared the placement.
// Keep stopping the original owner instead of replacing the operator with a
// coordinator-local stop carrying a newer epoch.
repeatedOp := oc.StopChangefeedWithMaintainerEpoch(context.Background(), cfID, false, 20)
require.Same(t, op, repeatedOp)
repeatedReqMsg := repeatedOp.Schedule()
require.Equal(t, owner.ID, repeatedReqMsg.To)
repeatedReq := repeatedReqMsg.Message[0].(*heartbeatpb.RemoveMaintainerRequest)
require.Equal(t, uint64(10), repeatedReq.MaintainerEpoch)
require.False(t, repeatedReq.Removed)
}

func TestController_StopRemoteMaintainerWithMaintainerEpoch(t *testing.T) {
Expand Down Expand Up @@ -143,6 +158,37 @@ func TestController_StopRemoteMaintainerWithMaintainerEpoch(t *testing.T) {
require.True(t, op.IsFinished())
}

func TestController_StopChangefeedDoesNotReuseStaleOwnerCleanup(t *testing.T) {
changefeedDB := changefeed.NewChangefeedDB(1216)
ctrl := gomock.NewController(t)
backend := mock_changefeed.NewMockBackend(ctrl)
oc, _, nodeManager := newOperatorControllerForTest(t, changefeedDB, backend, nil)
staleOwner := node.NewInfo("localhost:8301", "")
currentOwner := node.NewInfo("localhost:8302", "")
nodeManager.GetAliveNodes()[staleOwner.ID] = staleOwner
nodeManager.GetAliveNodes()[currentOwner.ID] = currentOwner

cfID := common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName)
cf := changefeed.NewChangefeed(cfID, &config.ChangeFeedInfo{
ChangefeedID: cfID,
Config: config.GetDefaultReplicaConfig(),
SinkURI: "mysql://127.0.0.1:3306",
Epoch: 20,
}, 1, true)
changefeedDB.AddReplicatingMaintainer(cf, currentOwner.ID)

staleOp := oc.StopRemoteMaintainerWithMaintainerEpoch(cfID, staleOwner.ID, false, 10)
require.Equal(t, staleOwner.ID, staleOp.Schedule().To)

backend.EXPECT().SetChangefeedProgress(gomock.Any(), cfID, config.ProgressNone).Return(nil).Times(1)
currentOp := oc.StopChangefeedWithMaintainerEpoch(context.Background(), cfID, false, 20)
require.NotSame(t, staleOp, currentOp)
currentReqMsg := currentOp.Schedule()
require.Equal(t, currentOwner.ID, currentReqMsg.To)
currentReq := currentReqMsg.Message[0].(*heartbeatpb.RemoveMaintainerRequest)
require.Equal(t, uint64(20), currentReq.MaintainerEpoch)
}

func TestController_AddOperator(t *testing.T) {
changefeedDB := changefeed.NewChangefeedDB(1216)
ctrl := gomock.NewController(t)
Expand Down
10 changes: 10 additions & 0 deletions coordinator/operator/operator_stop.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,13 @@ import (
"go.uber.org/zap"
)

type stopChangefeedKind int

const (
stopChangefeedKindCurrentPlacement stopChangefeedKind = iota
stopChangefeedKindStaleOwner
)

// StopChangefeedOperator is an operator to remove a maintainer from a node
type StopChangefeedOperator struct {
keyspaceID uint32
Expand All @@ -38,6 +45,7 @@ type StopChangefeedOperator struct {
coordinatorNodeID node.ID
backend changefeed.Backend
maintainerEpoch uint64
kind stopChangefeedKind
}

func NewStopChangefeedOperator(
Expand All @@ -48,6 +56,7 @@ func NewStopChangefeedOperator(
backend changefeed.Backend,
removed bool,
maintainerEpoch uint64,
kind stopChangefeedKind,
) *StopChangefeedOperator {
return &StopChangefeedOperator{
keyspaceID: keyspaceID,
Expand All @@ -57,6 +66,7 @@ func NewStopChangefeedOperator(
coordinatorNodeID: coordinatorNode,
backend: backend,
maintainerEpoch: maintainerEpoch,
kind: kind,
}
}

Expand Down
8 changes: 4 additions & 4 deletions coordinator/operator/operator_stop_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ func TestStopChangefeedOperator_OnNodeRemove(t *testing.T) {

ctrl := gomock.NewController(t)
backend := mock_changefeed.NewMockBackend(ctrl)
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, true, 10)
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, true, 10, stopChangefeedKindCurrentPlacement)
op.OnNodeRemove("n1")
require.Equal(t, "n2", op.nodeID.String())
require.False(t, op.finished.Load())
Expand All @@ -54,7 +54,7 @@ func TestStopChangefeedOperator_OnTaskRemoved(t *testing.T) {
},
1, true)
changefeedDB.AddReplicatingMaintainer(cf, "n1")
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", nil, true, 10)
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", nil, true, 10, stopChangefeedKindCurrentPlacement)
op.OnTaskRemoved()
require.True(t, op.finished.Load())
}
Expand All @@ -72,11 +72,11 @@ func TestStopChangefeedOperator_PostFinish(t *testing.T) {

ctrl := gomock.NewController(t)
backend := mock_changefeed.NewMockBackend(ctrl)
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, true, 10)
op := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, true, 10, stopChangefeedKindCurrentPlacement)
backend.EXPECT().DeleteChangefeed(gomock.Any(), cfID).Return(errors.New("err"))
op.PostFinish()

op2 := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, false, 10)
op2 := NewStopChangefeedOperator(common.DefaultKeyspaceID, cfID, "n1", "n2", backend, false, 10, stopChangefeedKindCurrentPlacement)
backend.EXPECT().SetChangefeedProgress(gomock.Any(), cfID, config.ProgressNone).Return(errors.New("err"))
op2.PostFinish()
}
Loading