diff --git a/downstreamadapter/dispatcher/event_dispatcher.go b/downstreamadapter/dispatcher/event_dispatcher.go index 7d04db09c5..35d1ee1979 100644 --- a/downstreamadapter/dispatcher/event_dispatcher.go +++ b/downstreamadapter/dispatcher/event_dispatcher.go @@ -152,9 +152,10 @@ func (d *EventDispatcher) Remove() { d.removeDispatcher() } -// EmitBootstrap emits the table bootstrap event in a blocking way after changefeed started -// It will return after the bootstrap event is sent. -func (d *EventDispatcher) EmitBootstrap() bool { +// EmitBootstrap emits the table bootstrap event in a blocking way after changefeed started. +// It will return after the bootstrap event is sent, or when shouldStop asks it +// to stop because the local write path has been fenced. +func (d *EventDispatcher) EmitBootstrap(shouldStop func() bool) bool { bootstrap := loadBootstrapState(&d.BootstrapState) switch bootstrap { case BootstrapFinished: @@ -204,9 +205,15 @@ func (d *EventDispatcher) EmitBootstrap() bool { if table.IsView() { continue } + if shouldStop() { + return false + } ddlEvent := codec.NewBootstrapDDLEvent(table) err := d.sink.WriteBlockEvent(ddlEvent) if err != nil { + if shouldStop() { + return false + } log.Error("send bootstrap message failed", zap.Stringer("changefeed", d.sharedInfo.changefeedID), zap.Int("tables", len(currentTables)), @@ -216,6 +223,9 @@ func (d *EventDispatcher) EmitBootstrap() bool { d.HandleError(errors.ErrExecDDLFailed.GenWithStackByArgs()) return true } + if shouldStop() { + return false + } } storeBootstrapState(&d.BootstrapState, BootstrapFinished) log.Info("send bootstrap messages finished", diff --git a/downstreamadapter/dispatchermanager/dispatcher_manager.go b/downstreamadapter/dispatchermanager/dispatcher_manager.go index 434772adcf..7ec156ea55 100644 --- a/downstreamadapter/dispatchermanager/dispatcher_manager.go +++ b/downstreamadapter/dispatchermanager/dispatcher_manager.go @@ -54,6 +54,21 @@ const ( blockStatusBufferSize = 16 * 1024 ) +// IsWritePathClosedError reports whether err means the local write path has +// already been fenced. Callers should stop the in-flight local request instead +// of treating it as a successful dispatcher creation. +func IsWritePathClosedError(err error) bool { + if err == nil { + return false + } + code, ok := errors.RFCCode(err) + return ok && code == errors.ErrDispatcherManagerWritePathClosed.RFCCode() +} + +func newWritePathClosedError() error { + return errors.ErrDispatcherManagerWritePathClosed.FastGenByArgs() +} + /* DispatcherManager manages dispatchers for a changefeed instance with responsibilities including: @@ -140,8 +155,10 @@ type DispatcherManager struct { latestWatermark Watermark latestRedoWatermark Watermark - closing atomic.Bool - closed atomic.Bool + closing atomic.Bool + closed atomic.Bool + writePathMu sync.Mutex + writePathClosed atomic.Bool // removeChangefeedRequested is sticky once any close request asks for removed=true. // A later removed=false request must not downgrade the final cleanup semantics. removeChangefeedRequested atomic.Bool @@ -191,7 +208,8 @@ func NewDispatcherManager( startTs uint64, maintainerID node.ID, newChangefeed bool, -) (*DispatcherManager, error) { + registerInitializing func(*DispatcherManager) bool, +) (manager *DispatcherManager, err error) { failpoint.Inject("NewDispatcherManagerDelay", nil) ctx, cancel := context.WithCancel(context.Background()) @@ -207,7 +225,7 @@ func NewDispatcherManager( integrityCfg = cfConfig.SinkConfig.Integrity.ToPB() } - manager := &DispatcherManager{ + manager = &DispatcherManager{ ctx: ctx, dispatcherMap: newDispatcherMap[*dispatcher.EventDispatcher](), currentOperatorMap: sync.Map{}, @@ -239,6 +257,17 @@ func NewDispatcherManager( // Set the epoch and maintainerID of the event dispatcher manager manager.meta.maintainerEpoch = cfConfig.Epoch manager.meta.maintainerID = maintainerID + cleanupManager := manager + defer func() { + if err != nil && cleanupManager != nil { + cleanupManager.LocalFence() + manager = nil + } + }() + // The manager must be fenceable before any write-capable resource is initialized. + if registerInitializing != nil && !registerInitializing(manager) { + return nil, newWritePathClosedError() + } // Set Sync Point Config var syncPointConfig *syncpoint.SyncPointConfig @@ -250,11 +279,18 @@ func NewDispatcherManager( } } - var err error - manager.sink, err = sink.New(ctx, manager.config, manager.changefeedID) + createdSink, err := sink.New(ctx, manager.config, manager.changefeedID) if err != nil { return nil, errors.Trace(err) } + manager.writePathMu.Lock() + if manager.writePathClosed.Load() { + manager.writePathMu.Unlock() + createdSink.Close() + return nil, newWritePathClosedError() + } + manager.sink = createdSink + manager.writePathMu.Unlock() sinkType := manager.sink.SinkType() if sinkType != common.KafkaSinkType { @@ -285,7 +321,7 @@ func NewDispatcherManager( return nil, err } // Create shared info for all dispatchers - manager.sharedInfo = dispatcher.NewSharedInfo( + sharedInfo := dispatcher.NewSharedInfo( manager.changefeedID, manager.config.TimeZone, manager.config.BDRMode, @@ -300,10 +336,24 @@ func NewDispatcherManager( blockStatusBufferSize, make(chan error, 1), ) + manager.writePathMu.Lock() + if manager.writePathClosed.Load() { + manager.writePathMu.Unlock() + sharedInfo.Close() + return nil, newWritePathClosedError() + } + manager.sharedInfo = sharedInfo + manager.writePathMu.Unlock() // Register Event Dispatcher Manager in HeartBeatCollector, // which is responsible for communication with the maintainer. + manager.writePathMu.Lock() + if manager.writePathClosed.Load() { + manager.writePathMu.Unlock() + return nil, newWritePathClosedError() + } err = appcontext.GetService[*HeartBeatCollector](appcontext.HeartbeatCollector).RegisterDispatcherManager(manager) + manager.writePathMu.Unlock() if err != nil { return nil, errors.Trace(err) } @@ -390,33 +440,61 @@ func (e *DispatcherManager) NewTableTriggerEventDispatcher(id *heartbeatpb.Dispa if err != nil { return errors.Trace(err) } + tableTriggerDispatcher := e.GetTableTriggerEventDispatcher() + if tableTriggerDispatcher == nil { + if e.writePathClosed.Load() { + return newWritePathClosedError() + } + return errors.ErrChangefeedInitTableTriggerDispatcherFailed. + FastGenByArgs("table trigger event dispatcher was not created") + } log.Info("table trigger event dispatcher created", zap.Stringer("changefeedID", e.changefeedID), - zap.Stringer("dispatcherID", e.GetTableTriggerEventDispatcher().GetId()), - zap.Uint64("startTs", e.GetTableTriggerEventDispatcher().GetStartTs()), + zap.Stringer("dispatcherID", tableTriggerDispatcher.GetId()), + zap.Uint64("startTs", tableTriggerDispatcher.GetStartTs()), ) return nil } func (e *DispatcherManager) InitalizeTableTriggerEventDispatcher(schemaInfo []*heartbeatpb.SchemaInfo) error { - if e.GetTableTriggerEventDispatcher() == nil { + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + return newWritePathClosedError() + } + tableTriggerDispatcher := e.GetTableTriggerEventDispatcher() + if tableTriggerDispatcher == nil { + e.writePathMu.Unlock() return nil } - needAddDispatcher, err := e.GetTableTriggerEventDispatcher().InitializeTableSchemaStore(schemaInfo) + needAddDispatcher, err := tableTriggerDispatcher.InitializeTableSchemaStore(schemaInfo) if err != nil { + e.writePathMu.Unlock() return errors.Trace(err) } + e.writePathMu.Unlock() if !needAddDispatcher { return nil } + if e.writePathClosed.Load() { + return newWritePathClosedError() + } // before bootstrap finished, cannot send any event. - success := e.GetTableTriggerEventDispatcher().EmitBootstrap() + success := tableTriggerDispatcher.EmitBootstrap(e.writePathClosed.Load) + if e.writePathClosed.Load() { + return newWritePathClosedError() + } if !success { return errors.ErrDispatcherFailed.GenWithStackByArgs() } + e.writePathMu.Lock() + defer e.writePathMu.Unlock() + if e.writePathClosed.Load() { + return newWritePathClosedError() + } // table trigger event dispatcher can register to event collector to receive events after finish the initial table schema store from the maintainer. - appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector).AddDispatcher(e.GetTableTriggerEventDispatcher(), e.sinkQuota) + appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector).AddDispatcher(tableTriggerDispatcher, e.sinkQuota) // when sink is not mysql-class, table trigger event dispatcher need to receive the checkpointTs message from maintainer. if e.sink.SinkType() != common.MysqlSinkType { @@ -455,6 +533,9 @@ func (e *DispatcherManager) getTableRecoveryInfoFromMysqlSink(tableIds, startTsL // 1. newEventDispatchers is called by NewTableTriggerEventDispatcher(just means when creating table trigger event dispatcher) // 2. changefeed is total new created, or resumed with overwriteCheckpointTs func (e *DispatcherManager) newEventDispatchers(infos map[common.DispatcherID]dispatcherCreateInfo, removeDDLTs bool) error { + if e.writePathClosed.Load() { + return newWritePathClosedError() + } start := time.Now() currentPdTs := e.pdClock.CurrentTS() @@ -518,6 +599,12 @@ func (e *DispatcherManager) newEventDispatchers(infos map[common.DispatcherID]di e.heartBeatTask = newHeartBeatTask(e) } + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + d.Remove() + return newWritePathClosedError() + } if d.IsTableTriggerDispatcher() { if util.GetOrZero(e.config.SinkConfig.SendAllBootstrapAtStart) { d.BootstrapState = dispatcher.BootstrapNotStarted @@ -533,6 +620,7 @@ func (e *DispatcherManager) newEventDispatchers(infos map[common.DispatcherID]di seq := e.dispatcherMap.Set(id, d) d.SetSeq(seq) + e.writePathMu.Unlock() if d.IsTableTriggerDispatcher() { e.metricTableTriggerEventDispatcherCount.Inc() @@ -900,7 +988,14 @@ func (e *DispatcherManager) mergeEventDispatcher(dispatcherIDs []common.Dispatch zap.Stringer("dispatcherID", mergedDispatcherID), zap.String("tableSpan", common.FormatTableSpan(mergedSpan))) + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + mergedDispatcher.Remove() + return nil + } registerMergeDispatcher(e.changefeedID, dispatcherIDs, e.dispatcherMap, mergedDispatcherID, mergedDispatcher, e.schemaIDToDispatchers, e.metricEventDispatcherCount, e.sinkQuota) + e.writePathMu.Unlock() return newMergeCheckTask(e, mergedDispatcher, dispatcherIDs) } @@ -914,45 +1009,88 @@ func (e *DispatcherManager) TryClose(removeChangefeed bool) bool { e.tryScheduleRemoveChangefeedCleanup() return true } - if e.closing.Load() { + if !e.closing.CompareAndSwap(false, true) { return e.closed.Load() } - e.closing.Store(true) go e.close() return false } +// LocalFence stops the local write path immediately without waiting for +// dispatcher progress to drain. The remaining cleanup continues asynchronously. +func (e *DispatcherManager) LocalFence() { + if e.closed.Load() { + return + } + startClose := e.closing.CompareAndSwap(false, true) + e.stopWritePath(true) + if startClose { + go e.finishClose() + } +} + func (e *DispatcherManager) close() { log.Info("closing event dispatcher manager", zap.Stringer("changefeedID", e.changefeedID)) - defer e.closing.Store(false) + e.stopWritePath(false) + e.finishClose() +} + +func (e *DispatcherManager) stopWritePath(cancelFirst bool) { + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + return + } + e.writePathClosed.Store(true) + e.writePathMu.Unlock() + + log.Info("stopping dispatcher manager write path", + zap.Stringer("changefeedID", e.changefeedID)) + + if cancelFirst && e.cancel != nil { + e.cancel() + } + if e.IsRedoEnabled() && e.redoSink != nil { closeAllDispatchers(e.changefeedID, e.redoDispatcherMap, e.redoSink.SinkType()) log.Info("closed all redo dispatchers", zap.Stringer("changefeedID", e.changefeedID)) - err := appcontext.GetService[*HeartBeatCollector](appcontext.HeartbeatCollector).RemoveRedoMessage(e.changefeedID) - if err != nil { - log.Error("remove redo message failed", + if heartbeatCollector, ok := appcontext.TryGetService[*HeartBeatCollector](appcontext.HeartbeatCollector); ok { + err := heartbeatCollector.RemoveRedoMessage(e.changefeedID) + if err != nil { + log.Error("remove redo message failed", + zap.Stringer("changefeedID", e.changefeedID), + zap.Error(err), + ) + } + } else { + log.Warn("heartbeat collector is not available when stopping redo write path", zap.Stringer("changefeedID", e.changefeedID), - zap.Error(err), ) - return } } - closeAllDispatchers(e.changefeedID, e.dispatcherMap, e.sink.SinkType()) + if e.sink != nil { + closeAllDispatchers(e.changefeedID, e.dispatcherMap, e.sink.SinkType()) + } log.Info("closed all event dispatchers", zap.Stringer("changefeedID", e.changefeedID)) - err := appcontext.GetService[*HeartBeatCollector](appcontext.HeartbeatCollector).RemoveDispatcherManager(e.changefeedID) - if err != nil { - log.Error("remove dispatcher manager from heartbeat collector failed", + if heartbeatCollector, ok := appcontext.TryGetService[*HeartBeatCollector](appcontext.HeartbeatCollector); ok { + err := heartbeatCollector.RemoveDispatcherManager(e.changefeedID) + if err != nil { + log.Error("remove dispatcher manager from heartbeat collector failed", + zap.Stringer("changefeedID", e.changefeedID), + zap.Error(err), + ) + } + } else { + log.Warn("heartbeat collector is not available when stopping dispatcher manager write path", zap.Stringer("changefeedID", e.changefeedID), - zap.Error(err), ) - return } // heartbeatTask only will be generated when create new dispatchers. @@ -963,10 +1101,9 @@ func (e *DispatcherManager) close() { e.heartBeatTask.Cancel() } - // Cancel the context to signal all dependent components to stop. - // This is important to prevent `e.sink.Close() / e.sharedInfo.Close()` from blocking, - // especially when a long-running DDL is being executed by the sink. - e.cancel() + if !cancelFirst && e.cancel != nil { + e.cancel() + } if e.sharedInfo != nil { e.sharedInfo.Close() @@ -977,10 +1114,31 @@ func (e *DispatcherManager) close() { if e.IsRedoEnabled() && e.redoSink != nil { e.redoSink.Close() } - e.sink.Close() + if e.sink != nil { + e.sink.Close() + } log.Info("sink closed", zap.Stringer("changefeedID", e.changefeedID)) +} +func (e *DispatcherManager) addCheckpointTs(checkpointTs uint64) { + if e.writePathClosed.Load() { + return + } + if e.GetTableTriggerEventDispatcher() == nil || e.sink == nil { + return + } + if e.writePathClosed.Load() { + return + } + e.sink.AddCheckpointTs(checkpointTs) +} + +func (e *DispatcherManager) finishClose() { + defer e.closing.Store(false) e.wg.Wait() + if !e.closed.CompareAndSwap(false, true) { + return + } e.removeTaskHandles.Range(func(key, value interface{}) bool { handle := value.(*threadpool.TaskHandle) @@ -990,7 +1148,6 @@ func (e *DispatcherManager) close() { e.cleanMetrics() - e.closed.Store(true) e.tryScheduleRemoveChangefeedCleanup() log.Info("event dispatcher manager closed", zap.Stringer("changefeedID", e.changefeedID)) diff --git a/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go b/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go index db69c69f80..5550497b43 100644 --- a/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go +++ b/downstreamadapter/dispatchermanager/dispatcher_manager_redo.go @@ -52,21 +52,36 @@ func initRedoComponet( if !manager.IsRedoEnabled() { return nil } + if manager.writePathClosed.Load() { + return newWritePathClosedError() + } var err error - manager.redoDispatcherMap = newDispatcherMap[*dispatcher.RedoDispatcher]() - manager.redoSink, err = redo.New(ctx, changefeedID, manager.config.Consistent) + redoDispatcherMap := newDispatcherMap[*dispatcher.RedoDispatcher]() + redoSink, err := redo.New(ctx, changefeedID, manager.config.Consistent) if err != nil { return err } - manager.redoSchemaIDToDispatchers = dispatcher.NewSchemaIDToDispatchers() + redoSchemaIDToDispatchers := dispatcher.NewSchemaIDToDispatchers() totalQuota := manager.sinkQuota consistentMemoryUsage := manager.config.Consistent.MemoryUsage if consistentMemoryUsage == nil { consistentMemoryUsage = config.GetDefaultReplicaConfig().Consistent.MemoryUsage } - manager.redoQuota = totalQuota * consistentMemoryUsage.MemoryQuotaPercentage / 100 - manager.sinkQuota = totalQuota - manager.redoQuota + redoQuota := totalQuota * consistentMemoryUsage.MemoryQuotaPercentage / 100 + + manager.writePathMu.Lock() + if manager.writePathClosed.Load() { + manager.writePathMu.Unlock() + redoSink.Close() + return newWritePathClosedError() + } + manager.redoDispatcherMap = redoDispatcherMap + manager.redoSink = redoSink + manager.redoSchemaIDToDispatchers = redoSchemaIDToDispatchers + manager.redoQuota = redoQuota + manager.sinkQuota = totalQuota - redoQuota + manager.writePathMu.Unlock() // init table trigger redo dispatcher when tableTriggerRedoDispatcherID is not nil if tableTriggerRedoDispatcherID != nil { @@ -81,7 +96,16 @@ func initRedoComponet( manager.metricRedoCreateDispatcherDuration = metrics.CreateDispatcherDuration.WithLabelValues(changefeedID.Keyspace(), changefeedID.Name(), "redoDispatcher") // RedoMessageDs need register on every node + manager.writePathMu.Lock() + if manager.writePathClosed.Load() { + manager.writePathMu.Unlock() + return newWritePathClosedError() + } appcontext.GetService[*HeartBeatCollector](appcontext.HeartbeatCollector).RegisterRedoMessageDs(manager) + manager.writePathMu.Unlock() + if manager.writePathClosed.Load() { + return newWritePathClosedError() + } manager.wg.Add(1) go func() { defer manager.wg.Done() @@ -113,6 +137,13 @@ func (e *DispatcherManager) NewTableTriggerRedoDispatcher(id *heartbeatpb.Dispat // redo meta should keep the same node with table trigger event dispatcher // table trigger event dispatcher and table trigger redo dispatcher must exist on the same node redoDispatcher := e.GetTableTriggerRedoDispatcher() + if redoDispatcher == nil { + if e.writePathClosed.Load() { + return newWritePathClosedError() + } + return errors.ErrChangefeedInitTableTriggerDispatcherFailed. + FastGenByArgs("table trigger redo dispatcher was not created") + } redoDispatcher.SetRedoMeta(e.ctx, e.config.Consistent) e.wg.Add(1) go func() { @@ -129,6 +160,9 @@ func (e *DispatcherManager) NewTableTriggerRedoDispatcher(id *heartbeatpb.Dispat } func (e *DispatcherManager) newRedoDispatchers(infos map[common.DispatcherID]dispatcherCreateInfo, removeDDLTs bool) error { + if e.writePathClosed.Load() { + return newWritePathClosedError() + } start := time.Now() dispatcherIds, _, startTsList, tableSpans, schemaIds, scheduleSkipDMLAsStartTsList := prepareCreateDispatcher(infos, e.redoDispatcherMap) @@ -166,6 +200,12 @@ func (e *DispatcherManager) newRedoDispatchers(infos map[common.DispatcherID]dis e.heartBeatTask = newHeartBeatTask(e) } + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + rd.Remove() + return newWritePathClosedError() + } if rd.IsTableTriggerDispatcher() { e.SetTableTriggerRedoDispatcher(rd) } else { @@ -175,6 +215,7 @@ func (e *DispatcherManager) newRedoDispatchers(infos map[common.DispatcherID]dis redoSeq := e.redoDispatcherMap.Set(rd.GetId(), rd) rd.SetSeq(redoSeq) + e.writePathMu.Unlock() if rd.IsTableTriggerDispatcher() { e.metricTableTriggerRedoDispatcherCount.Inc() @@ -232,7 +273,14 @@ func (e *DispatcherManager) mergeRedoDispatcher(dispatcherIDs []common.Dispatche zap.Stringer("dispatcherID", mergedDispatcherID), zap.String("tableSpan", common.FormatTableSpan(mergedSpan))) + e.writePathMu.Lock() + if e.writePathClosed.Load() { + e.writePathMu.Unlock() + mergedDispatcher.Remove() + return nil + } registerMergeDispatcher(e.changefeedID, dispatcherIDs, e.redoDispatcherMap, mergedDispatcherID, mergedDispatcher, e.redoSchemaIDToDispatchers, e.metricRedoEventDispatcherCount, e.redoQuota) + e.writePathMu.Unlock() return newMergeCheckTask(e, mergedDispatcher, dispatcherIDs) } @@ -273,21 +321,30 @@ func (e *DispatcherManager) closeRedoMeta(removeChangefeed bool) error { } func (e *DispatcherManager) InitalizeTableTriggerRedoDispatcher(schemaInfo []*heartbeatpb.SchemaInfo) error { - if e.GetTableTriggerRedoDispatcher() == nil { + e.writePathMu.Lock() + defer e.writePathMu.Unlock() + if e.writePathClosed.Load() { + return newWritePathClosedError() + } + tableTriggerRedoDispatcher := e.GetTableTriggerRedoDispatcher() + if tableTriggerRedoDispatcher == nil { return nil } - needAddDispatcher, err := e.GetTableTriggerRedoDispatcher().InitializeTableSchemaStore(schemaInfo) + needAddDispatcher, err := tableTriggerRedoDispatcher.InitializeTableSchemaStore(schemaInfo) if err != nil { return errors.Trace(err) } if !needAddDispatcher { return nil } - appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector).AddDispatcher(e.GetTableTriggerRedoDispatcher(), e.redoQuota) + appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector).AddDispatcher(tableTriggerRedoDispatcher, e.redoQuota) return nil } func (e *DispatcherManager) UpdateRedoMeta(checkpointTs, resolvedTs uint64) { + if e.writePathClosed.Load() { + return + } // only update meta on the one node d := e.GetTableTriggerRedoDispatcher() if d == nil { diff --git a/downstreamadapter/dispatchermanager/dispatcher_manager_test.go b/downstreamadapter/dispatchermanager/dispatcher_manager_test.go index 488f7f0230..98dd3101d8 100644 --- a/downstreamadapter/dispatchermanager/dispatcher_manager_test.go +++ b/downstreamadapter/dispatchermanager/dispatcher_manager_test.go @@ -26,10 +26,12 @@ import ( "github.com/pingcap/ticdc/downstreamadapter/sink/mock" "github.com/pingcap/ticdc/downstreamadapter/sink/mysql" "github.com/pingcap/ticdc/heartbeatpb" + "github.com/pingcap/ticdc/logservice/schemastore" "github.com/pingcap/ticdc/pkg/common" appcontext "github.com/pingcap/ticdc/pkg/common/context" "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/config" + "github.com/pingcap/ticdc/pkg/filter" "github.com/pingcap/ticdc/pkg/messaging" "github.com/pingcap/ticdc/pkg/metrics" "github.com/pingcap/ticdc/pkg/node" @@ -188,6 +190,86 @@ func TestCountIgnoreUpdateOnlyColumnsRules(t *testing.T) { } } +type bootstrapSchemaStoreForTest struct{} + +func (s *bootstrapSchemaStoreForTest) Name() string { return "bootstrap-schema-store-for-test" } + +func (s *bootstrapSchemaStoreForTest) Run(ctx context.Context) error { return nil } + +func (s *bootstrapSchemaStoreForTest) Close(ctx context.Context) error { return nil } + +func (s *bootstrapSchemaStoreForTest) GetAllPhysicalTables( + keyspaceMeta common.KeyspaceMeta, + snapTs uint64, + filter filter.Filter, +) ([]event.Table, error) { + return nil, nil +} + +func (s *bootstrapSchemaStoreForTest) RegisterTable( + keyspaceMeta common.KeyspaceMeta, + tableID int64, + startTs uint64, +) error { + return nil +} + +func (s *bootstrapSchemaStoreForTest) UnregisterTable( + keyspaceMeta common.KeyspaceMeta, + tableID int64, +) error { + return nil +} + +func (s *bootstrapSchemaStoreForTest) GetTableInfo( + keyspaceMeta common.KeyspaceMeta, + tableID int64, + ts uint64, +) (*common.TableInfo, error) { + return &common.TableInfo{ + TableName: common.TableName{ + Schema: "test", + Table: "t", + TableID: tableID, + }, + }, nil +} + +func (s *bootstrapSchemaStoreForTest) GetTableDDLEventState( + keyspaceMeta common.KeyspaceMeta, + tableID int64, +) (schemastore.DDLEventState, error) { + return schemastore.DDLEventState{}, nil +} + +func (s *bootstrapSchemaStoreForTest) FetchTableDDLEvents( + keyspaceMeta common.KeyspaceMeta, + dispatcherID common.DispatcherID, + tableID int64, + tableFilter filter.Filter, + start uint64, + end uint64, +) ([]event.DDLEvent, error) { + return nil, nil +} + +func (s *bootstrapSchemaStoreForTest) FetchTableTriggerDDLEvents( + keyspaceMeta common.KeyspaceMeta, + dispatcherID common.DispatcherID, + tableFilter filter.Filter, + start uint64, + limit int, +) ([]event.DDLEvent, uint64, error) { + return nil, 0, nil +} + +func (s *bootstrapSchemaStoreForTest) RegisterKeyspace( + ctx context.Context, + keyspaceMeta common.KeyspaceMeta, +) error { + return nil +} + func TestCollectComponentStatusWhenChangedWatermarkSeqNoFallback(t *testing.T) { manager := createTestManager(t) @@ -441,6 +523,306 @@ func TestTryCloseRemovedRequestAfterClosedReturnsImmediatelyAndTriggersCleanup(t require.True(t, manager.TryClose(true)) } +func TestLocalFenceCancelsWritePathWithoutWaitingForCleanup(t *testing.T) { + manager := createTestManager(t) + ctx, cancel := context.WithCancel(context.Background()) + manager.ctx = ctx + manager.cancel = cancel + + manager.wg.Add(1) + done := make(chan struct{}) + go func() { + manager.LocalFence() + close(done) + }() + + select { + case <-done: + case <-time.After(time.Second): + require.FailNow(t, "local fence should not wait for dispatcher manager cleanup") + } + require.ErrorIs(t, ctx.Err(), context.Canceled) + + manager.wg.Done() + require.Eventually(t, func() bool { + return manager.TryClose(false) + }, time.Second, 10*time.Millisecond) +} + +func TestLocalFenceDoesNotWaitForBootstrapWriteBlockEvent(t *testing.T) { + manager := createTestManager(t) + appcontext.SetService(appcontext.SchemaStore, &bootstrapSchemaStoreForTest{}) + heartbeatCollector := &HeartBeatCollector{} + heartbeatCollector.isClosed.Store(true) + appcontext.SetService(appcontext.HeartbeatCollector, heartbeatCollector) + + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + writeStarted := make(chan struct{}) + releaseWrite := make(chan struct{}) + mockSink := mock.NewMockSink(ctrl) + mockSink.EXPECT().SinkType().Return(common.KafkaSinkType).AnyTimes() + mockSink.EXPECT().IsNormal().Return(true).AnyTimes() + mockSink.EXPECT().SetTableSchemaStore(gomock.Any()).AnyTimes() + mockSink.EXPECT().Close().AnyTimes() + mockSink.EXPECT().WriteBlockEvent(gomock.Any()).DoAndReturn(func(blockEvent event.BlockEvent) error { + close(writeStarted) + <-releaseWrite + blockEvent.PostFlush() + return nil + }).Times(1) + manager.sink = mockSink + + dispatcherID := common.NewDispatcherID() + var redoTs atomic.Uint64 + redoTs.Store(math.MaxUint64) + tableTriggerDispatcher := dispatcher.NewEventDispatcher( + dispatcherID, + common.KeyspaceDDLSpan(common.DefaultKeyspaceID), + 1, + 0, + manager.schemaIDToDispatchers, + false, + false, + 0, + manager.sink, + manager.sharedInfo, + false, + &redoTs, + ) + tableTriggerDispatcher.BootstrapState = dispatcher.BootstrapNotStarted + manager.SetTableTriggerEventDispatcher(tableTriggerDispatcher) + manager.dispatcherMap.Set(dispatcherID, tableTriggerDispatcher) + + initErrCh := make(chan error, 1) + go func() { + initErrCh <- manager.InitalizeTableTriggerEventDispatcher([]*heartbeatpb.SchemaInfo{ + { + SchemaID: 1, + SchemaName: "test", + Tables: []*heartbeatpb.TableInfo{ + {TableID: 11, TableName: "t"}, + }, + }, + }) + }() + + select { + case <-writeStarted: + case <-time.After(time.Second): + require.FailNow(t, "bootstrap should reach WriteBlockEvent") + } + + fenceDone := make(chan struct{}) + go func() { + manager.LocalFence() + close(fenceDone) + }() + select { + case <-fenceDone: + case <-time.After(100 * time.Millisecond): + require.FailNow(t, "local fence should not wait for blocked bootstrap WriteBlockEvent") + } + + close(releaseWrite) + select { + case err := <-initErrCh: + require.True(t, IsWritePathClosedError(err)) + case <-time.After(time.Second): + require.FailNow(t, "bootstrap initialization should return after write unblocks") + } + require.False(t, appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector).HasDispatcher(dispatcherID)) +} + +func TestNewDispatcherManagerReturnsFenceErrorWhenInitializingRegistrationRejected(t *testing.T) { + appcontext.SetService(appcontext.DefaultPDClock, pdutil.NewClock4Test()) + + replicaConfig := config.GetDefaultReplicaConfig() + cfConfig := &config.ChangefeedConfig{ + SinkURI: "blackhole://", + SinkConfig: replicaConfig.Sink, + Filter: replicaConfig.Filter, + MemoryQuota: util.GetOrZero(replicaConfig.MemoryQuota), + TimeZone: "system", + Consistent: replicaConfig.Consistent, + } + changefeedID := common.NewChangeFeedIDWithName("test", common.DefaultKeyspaceName) + var hookCalled atomic.Bool + var initializingManager *DispatcherManager + + manager, err := NewDispatcherManager( + common.DefaultKeyspaceID, + changefeedID, + cfConfig, + nil, + nil, + 1, + node.ID("maintainer"), + true, + func(manager *DispatcherManager) bool { + hookCalled.Store(true) + initializingManager = manager + return false + }, + ) + + require.Nil(t, manager) + require.True(t, hookCalled.Load()) + require.NotNil(t, initializingManager) + require.True(t, IsWritePathClosedError(err)) + require.True(t, initializingManager.writePathClosed.Load()) + require.Eventually(t, func() bool { + return initializingManager.TryClose(false) + }, time.Second, 10*time.Millisecond) +} + +func TestCheckpointTsMessageHandlerSkipsWriteAfterLocalFence(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockSink := mock.NewMockSink(ctrl) + mockSink.EXPECT().SinkType().Return(common.MysqlSinkType).AnyTimes() + mockSink.EXPECT().Close().AnyTimes() + mockSink.EXPECT().AddCheckpointTs(gomock.Any()).Times(0) + + manager := createTestManager(t) + manager.sink = mockSink + manager.SetTableTriggerEventDispatcher(&dispatcher.EventDispatcher{}) + heartbeatCollector := &HeartBeatCollector{} + heartbeatCollector.isClosed.Store(true) + appcontext.SetService(appcontext.HeartbeatCollector, heartbeatCollector) + + handlerStarted := make(chan struct{}) + proceed := make(chan struct{}) + handlerDone := make(chan struct{}) + go func() { + close(handlerStarted) + <-proceed + handler := &CheckpointTsMessageHandler{} + handler.Handle(manager, NewCheckpointTsMessage(&heartbeatpb.CheckpointTsMessage{ + ChangefeedID: manager.changefeedID.ToPB(), + CheckpointTs: 100, + })) + close(handlerDone) + }() + + select { + case <-handlerStarted: + case <-time.After(time.Second): + require.FailNow(t, "handler goroutine should start") + } + manager.LocalFence() + close(proceed) + select { + case <-handlerDone: + case <-time.After(time.Second): + require.FailNow(t, "handler should return without writing checkpoint ts") + } +} + +func TestLocalFenceWithRedoEnabledBeforeRedoSinkInitialized(t *testing.T) { + manager := createTestManager(t) + manager.redoEnabled = true + manager.redoSink = nil + + require.NotPanics(t, func() { + manager.LocalFence() + }) + require.Eventually(t, func() bool { + return manager.TryClose(false) + }, time.Second, 10*time.Millisecond) +} + +func TestNewTableTriggerDispatchersReturnFenceErrorWhenWritePathClosed(t *testing.T) { + manager := createTestManager(t) + manager.writePathClosed.Store(true) + + dispatcherID := common.NewDispatcherID() + require.NotPanics(t, func() { + err := manager.NewTableTriggerEventDispatcher(dispatcherID.ToPB(), 1, false) + require.True(t, IsWritePathClosedError(err)) + }) + require.Nil(t, manager.GetTableTriggerEventDispatcher()) + + redoDispatcherID := common.NewDispatcherID() + require.NotPanics(t, func() { + err := manager.NewTableTriggerRedoDispatcher(redoDispatcherID.ToPB(), 1, false) + require.True(t, IsWritePathClosedError(err)) + }) + require.Nil(t, manager.GetTableTriggerRedoDispatcher()) +} + +func TestInitializeTableTriggerEventDispatcherReturnsFenceErrorWhenWritePathClosed(t *testing.T) { + manager := createTestManager(t) + dispatcherID := common.NewDispatcherID() + var redoTs atomic.Uint64 + redoTs.Store(math.MaxUint64) + tableTriggerDispatcher := dispatcher.NewEventDispatcher( + dispatcherID, + common.KeyspaceDDLSpan(common.DefaultKeyspaceID), + 1, + 0, + manager.schemaIDToDispatchers, + false, + false, + 0, + manager.sink, + manager.sharedInfo, + false, + &redoTs, + ) + manager.SetTableTriggerEventDispatcher(tableTriggerDispatcher) + manager.dispatcherMap.Set(dispatcherID, tableTriggerDispatcher) + manager.writePathClosed.Store(true) + + err := manager.InitalizeTableTriggerEventDispatcher(nil) + require.True(t, IsWritePathClosedError(err)) + + err = manager.InitalizeTableTriggerRedoDispatcher(nil) + require.True(t, IsWritePathClosedError(err)) + + eventCollector := appcontext.GetService[*eventcollector.EventCollector](appcontext.EventCollector) + require.False(t, eventCollector.HasDispatcher(dispatcherID)) +} + +func TestCreateDispatcherByInfoKeepsCreateOperatorWhenFenced(t *testing.T) { + manager := createTestManager(t) + manager.writePathClosed.Store(true) + dispatcherID := common.NewDispatcherID() + createReq := NewSchedulerDispatcherRequest(&heartbeatpb.ScheduleDispatcherRequest{ + ChangefeedID: manager.changefeedID.ToPB(), + Config: &heartbeatpb.DispatcherConfig{ + DispatcherID: dispatcherID.ToPB(), + Span: &heartbeatpb.TableSpan{ + TableID: 1, + }, + StartTs: 1, + Mode: common.DefaultMode, + }, + ScheduleAction: heartbeatpb.ScheduleAction_Create, + OperatorType: heartbeatpb.OperatorType_O_Add, + }) + manager.currentOperatorMap.Store(dispatcherID, createReq) + + createDispatcherByInfo(manager, map[common.DispatcherID]dispatcherCreateInfo{ + dispatcherID: { + Id: dispatcherID, + TableSpan: &heartbeatpb.TableSpan{ + TableID: 1, + }, + StartTs: 1, + SchemaID: 1, + }, + }, nil) + + _, dispatcherExists := manager.dispatcherMap.Get(dispatcherID) + require.False(t, dispatcherExists) + operator, operatorExists := manager.currentOperatorMap.Load(dispatcherID) + require.True(t, operatorExists) + require.Equal(t, createReq, operator) +} + func TestMergeDispatcherExistingID(t *testing.T) { manager := createTestManager(t) diff --git a/downstreamadapter/dispatchermanager/helper.go b/downstreamadapter/dispatchermanager/helper.go index 68d1b3e22c..9f0ad5b5d6 100644 --- a/downstreamadapter/dispatchermanager/helper.go +++ b/downstreamadapter/dispatchermanager/helper.go @@ -409,40 +409,57 @@ func createDispatcherByInfo( if len(redoInfos) > 0 { err := dispatcherManager.newRedoDispatchers(redoInfos, false) if err != nil { - dispatcherManager.handleError(context.Background(), err) - } - for _, info := range redoInfos { - // Create requests are stored in currentOperatorMap before creation and should be deleted once the dispatcher is created. - if v, ok := dispatcherManager.currentOperatorMap.Load(info.Id); ok { - req := v.(SchedulerDispatcherRequest) - if req.ScheduleAction == heartbeatpb.ScheduleAction_Create { - log.Debug("delete current working add operator for redo dispatcher", - zap.String("changefeedID", dispatcherManager.changefeedID.String()), - zap.String("dispatcherID", info.Id.String()), - zap.Any("operator", req), - ) - dispatcherManager.currentOperatorMap.Delete(info.Id) - } + if IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed, keep add operators for redo dispatchers", + zap.String("changefeedID", dispatcherManager.changefeedID.String()), + zap.Int("count", len(redoInfos)), + zap.Error(err), + ) + return } + dispatcherManager.handleError(context.Background(), err) } + deleteCreatedOperators(dispatcherManager, redoInfos, dispatcherManager.redoDispatcherMap, "redo dispatcher") } if len(infos) > 0 { err := dispatcherManager.newEventDispatchers(infos, false) if err != nil { + if IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed, keep add operators", + zap.String("changefeedID", dispatcherManager.changefeedID.String()), + zap.Int("count", len(infos)), + zap.Error(err), + ) + return + } dispatcherManager.handleError(context.Background(), err) } - for _, info := range infos { - // Create requests are stored in currentOperatorMap before creation and should be deleted once the dispatcher is created. - if v, ok := dispatcherManager.currentOperatorMap.Load(info.Id); ok { - req := v.(SchedulerDispatcherRequest) - if req.ScheduleAction == heartbeatpb.ScheduleAction_Create { - log.Debug("delete current working add operator", - zap.String("changefeedID", dispatcherManager.changefeedID.String()), - zap.String("dispatcherID", info.Id.String()), - zap.Any("operator", req), - ) - dispatcherManager.currentOperatorMap.Delete(info.Id) - } + deleteCreatedOperators(dispatcherManager, infos, dispatcherManager.dispatcherMap, "dispatcher") + } +} + +func deleteCreatedOperators[T dispatcher.Dispatcher]( + dispatcherManager *DispatcherManager, + infos map[common.DispatcherID]dispatcherCreateInfo, + dispatcherMap *DispatcherMap[T], + dispatcherKind string, +) { + for _, info := range infos { + if _, exists := dispatcherMap.Get(info.Id); !exists { + continue + } + // Create requests are stored in currentOperatorMap before creation and + // should be deleted only after the dispatcher is actually created. + if v, ok := dispatcherManager.currentOperatorMap.Load(info.Id); ok { + req := v.(SchedulerDispatcherRequest) + if req.ScheduleAction == heartbeatpb.ScheduleAction_Create { + log.Debug("delete current working add operator", + zap.String("changefeedID", dispatcherManager.changefeedID.String()), + zap.String("dispatcherID", info.Id.String()), + zap.String("dispatcherKind", dispatcherKind), + zap.Any("operator", req), + ) + dispatcherManager.currentOperatorMap.Delete(info.Id) } } } @@ -603,7 +620,7 @@ func (h *CheckpointTsMessageHandler) Handle(dispatcherManager *DispatcherManager } if dispatcherManager.GetTableTriggerEventDispatcher() != nil { checkpointTsMessage := messages[0] - dispatcherManager.sink.AddCheckpointTs(checkpointTsMessage.CheckpointTs) + dispatcherManager.addCheckpointTs(checkpointTsMessage.CheckpointTs) } return false } diff --git a/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator.go b/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator.go index 49d3764955..04c41524ad 100644 --- a/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator.go +++ b/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator.go @@ -41,6 +41,9 @@ type DispatcherOrchestrator struct { mc messaging.MessageCenter mutex sync.Mutex // protect dispatcherManagers dispatcherManagers map[common.ChangeFeedID]*dispatchermanager.DispatcherManager + // initializingDispatcherManagers tracks managers that have been allocated + // but are not yet visible in dispatcherManagers. + initializingDispatcherManagers map[common.ChangeFeedID]*dispatchermanager.DispatcherManager // shards partition changefeed control messages by changefeed ID. Each shard keeps // the existing FIFO queue semantics, while different shards can process messages @@ -49,6 +52,9 @@ type DispatcherOrchestrator struct { // closed indicates Close has been invoked and no more messages should be enqueued. closed atomic.Bool + // fenced indicates this capture has lost local liveness and must stop local + // downstream writes without waiting for graceful dispatcher draining. + fenced atomic.Bool // msgGuardWaitGroup waits for in-flight RecvMaintainerRequest handlers before shutdown. msgGuardWaitGroup util.GuardedWaitGroup } @@ -62,9 +68,10 @@ const ( func New() *DispatcherOrchestrator { m := &DispatcherOrchestrator{ - mc: appcontext.GetService[messaging.MessageCenter](appcontext.MessageCenter), - dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), - shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), + mc: appcontext.GetService[messaging.MessageCenter](appcontext.MessageCenter), + dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + initializingDispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), } for i := range m.shards { m.shards[i] = newOrchestratorShard(m.processMessage) @@ -87,7 +94,7 @@ func (m *DispatcherOrchestrator) RecvMaintainerRequest( _ context.Context, msg *messaging.TargetMessage, ) error { - if !m.msgGuardWaitGroup.AddIf(func() bool { return !m.closed.Load() }) { + if !m.msgGuardWaitGroup.AddIf(func() bool { return !m.closed.Load() && !m.fenced.Load() }) { log.Debug("dispatcher orchestrator already closed, drop message", zap.Any("message", msg.Message)) return nil } @@ -139,6 +146,12 @@ func getPendingMessageKey(msg *messaging.TargetMessage) (pendingMessageKey, bool // processMessage dispatches a queued control message to the existing handler // implementation. Shards only change concurrency, not per-message behavior. func (m *DispatcherOrchestrator) processMessage(msg *messaging.TargetMessage) { + if m.fenced.Load() { + log.Info("dispatcher orchestrator is fenced, drop pending message", + zap.String("type", msg.Type.String())) + return + } + switch req := msg.Message[0].(type) { case *heartbeatpb.MaintainerBootstrapRequest: if err := m.handleBootstrapRequest(msg.From, req); err != nil { @@ -164,6 +177,9 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( from node.ID, req *heartbeatpb.MaintainerBootstrapRequest, ) error { + if m.fenced.Load() { + return nil + } cfId := common.NewChangefeedIDFromPB(req.ChangefeedID) cfConfig := &config.ChangefeedConfig{} @@ -182,6 +198,7 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( var err error if !exists { start := time.Now() + var initializingManager *dispatchermanager.DispatcherManager manager, err = dispatchermanager.NewDispatcherManager( req.KeyspaceId, cfId, @@ -191,9 +208,21 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( req.StartTs, from, req.IsNewChangefeed, + func(manager *dispatchermanager.DispatcherManager) bool { + initializingManager = manager + return m.registerInitializingDispatcherManager(cfId, manager) + }, ) + if initializingManager != nil { + m.removeInitializingDispatcherManager(cfId, initializingManager) + } // Fast return the error to maintainer. if err != nil { + if dispatchermanager.IsWritePathClosedError(err) || m.fenced.Load() || m.closed.Load() { + log.Info("dispatcher manager write path closed while creating dispatcher manager", + zap.Stringer("changefeedID", cfId), zap.Duration("duration", time.Since(start)), zap.Error(err)) + return nil + } log.Error("failed to create new dispatcher manager", zap.Any("changefeedID", cfId.Name()), zap.Duration("duration", time.Since(start)), zap.Error(err)) @@ -214,10 +243,18 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( return m.sendResponse(from, messaging.MaintainerManagerTopic, response) } m.mutex.Lock() + if m.fenced.Load() || m.closed.Load() { + m.mutex.Unlock() + manager.LocalFence() + return nil + } m.dispatcherManagers[cfId] = manager m.mutex.Unlock() metrics.DispatcherManagerGauge.WithLabelValues(cfId.Keyspace(), cfId.Name()).Inc() } else { + if m.fenced.Load() { + return nil + } // Check and potentially add a table trigger event dispatcher. // This is necessary during maintainer node migration, as the existing // dispatcher manager on the new node may not have a table trigger @@ -231,6 +268,11 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( false, ) if err != nil { + if dispatchermanager.IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed while creating table trigger event dispatcher", + zap.Stringer("changefeedID", cfId), zap.Error(err)) + return nil + } log.Error("failed to create new table trigger event dispatcher", zap.Stringer("changefeedID", cfId), zap.Error(err)) return m.handleDispatcherError(from, req.ChangefeedID, err) @@ -246,6 +288,11 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( false, ) if err != nil { + if dispatchermanager.IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed while creating table trigger redo dispatcher", + zap.Stringer("changefeedID", cfId), zap.Error(err)) + return nil + } log.Error("failed to create new table trigger redo dispatcher", zap.Stringer("changefeedID", cfId), zap.Error(err)) return m.handleDispatcherError(from, req.ChangefeedID, err) @@ -267,6 +314,11 @@ func (m *DispatcherOrchestrator) handleBootstrapRequest( zap.String("changefeed", cfId.Name()), zap.Uint64("epoch", cfConfig.Epoch)) } + if m.fenced.Load() { + manager.LocalFence() + return nil + } + var ( startTs uint64 redoStartTs uint64 @@ -291,6 +343,9 @@ func (m *DispatcherOrchestrator) handlePostBootstrapRequest( from node.ID, req *heartbeatpb.MaintainerPostBootstrapRequest, ) error { + if m.fenced.Load() { + return nil + } cfId := common.NewChangefeedIDFromPB(req.ChangefeedID) m.mutex.Lock() @@ -330,6 +385,11 @@ func (m *DispatcherOrchestrator) handlePostBootstrapRequest( // init table schema store err := manager.InitalizeTableTriggerEventDispatcher(req.Schemas) if err != nil { + if dispatchermanager.IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed while initializing table trigger event dispatcher", + zap.Any("changefeedID", cfId.Name()), zap.Error(err)) + return nil + } log.Error("failed to initialize table trigger event dispatcher", zap.Any("changefeedID", cfId.Name()), zap.Error(err)) return m.handleDispatcherError(from, req.ChangefeedID, err) @@ -337,12 +397,22 @@ func (m *DispatcherOrchestrator) handlePostBootstrapRequest( if manager.IsRedoReady() { err := manager.InitalizeTableTriggerRedoDispatcher(req.RedoSchemas) if err != nil { + if dispatchermanager.IsWritePathClosedError(err) { + log.Info("dispatcher manager write path closed while initializing table trigger redo dispatcher", + zap.Any("changefeedID", cfId.Name()), zap.Error(err)) + return nil + } log.Error("failed to initialize table trigger redo dispatcher", zap.Any("changefeedID", cfId.Name()), zap.Error(err)) return m.handleDispatcherError(from, req.ChangefeedID, err) } } + if m.fenced.Load() { + manager.LocalFence() + return nil + } + response := &heartbeatpb.MaintainerPostBootstrapResponse{ ChangefeedID: req.ChangefeedID, TableTriggerEventDispatcherId: req.TableTriggerEventDispatcherId, @@ -377,6 +447,70 @@ func (m *DispatcherOrchestrator) handleCloseRequest( return m.sendResponse(from, messaging.MaintainerTopic, response) } +// LocalFence stops all local dispatcher managers from writing downstream without +// waiting for dispatcher drain. Full resource cleanup is still handled by Close. +func (m *DispatcherOrchestrator) LocalFence() { + if !m.fenced.CompareAndSwap(false, true) { + m.localFenceManagers() + return + } + log.Warn("dispatcher orchestrator local fence triggered") + + m.mc.DeRegisterHandler(messaging.DispatcherManagerManagerTopic) + m.localFenceManagers() + m.msgGuardWaitGroup.Wait() + for _, shard := range m.shards { + shard.CloseAsync() + } + m.localFenceManagers() +} + +func (m *DispatcherOrchestrator) localFenceManagers() { + m.mutex.Lock() + managers := make([]*dispatchermanager.DispatcherManager, 0, + len(m.dispatcherManagers)+len(m.initializingDispatcherManagers)) + for _, manager := range m.dispatcherManagers { + managers = append(managers, manager) + } + for _, manager := range m.initializingDispatcherManagers { + managers = append(managers, manager) + } + m.mutex.Unlock() + + for _, manager := range managers { + manager.LocalFence() + } +} + +func (m *DispatcherOrchestrator) registerInitializingDispatcherManager( + cfID common.ChangeFeedID, + manager *dispatchermanager.DispatcherManager, +) bool { + m.mutex.Lock() + defer m.mutex.Unlock() + + if m.fenced.Load() || m.closed.Load() { + return false + } + if m.initializingDispatcherManagers == nil { + m.initializingDispatcherManagers = make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager) + } + m.initializingDispatcherManagers[cfID] = manager + return true +} + +func (m *DispatcherOrchestrator) removeInitializingDispatcherManager( + cfID common.ChangeFeedID, + manager *dispatchermanager.DispatcherManager, +) { + m.mutex.Lock() + defer m.mutex.Unlock() + + if m.initializingDispatcherManagers[cfID] == manager { + delete(m.initializingDispatcherManagers, cfID) + } +} + func createBootstrapResponse( changefeedID *heartbeatpb.ChangefeedID, manager *dispatchermanager.DispatcherManager, diff --git a/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator_test.go b/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator_test.go index ef88de6e57..0e3781249a 100644 --- a/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator_test.go +++ b/downstreamadapter/dispatcherorchestrator/dispatcher_orchestrator_test.go @@ -171,8 +171,9 @@ func newTestDispatcherOrchestrator() *DispatcherOrchestrator { // This test only exercises local routing through RecvMaintainerRequest, so it // needs shard state and the dispatcher manager map but not a message center. orchestrator := &DispatcherOrchestrator{ - dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), - shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), + dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + initializingDispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), } for i := range orchestrator.shards { orchestrator.shards[i] = newOrchestratorShard(func(msg *messaging.TargetMessage) {}) @@ -425,3 +426,92 @@ func TestGetPendingMessageKey_SupportedTypes(t *testing.T) { require.True(t, ok) require.Equal(t, pendingMessageKey{changefeedID: cfID, msgType: messaging.TypeMaintainerCloseRequest}, key) } + +func TestDispatcherOrchestratorLocalFenceDropsNewMessages(t *testing.T) { + mc, _, stop := messaging.NewMessageCenterForTest(t) + defer stop() + + orchestrator := &DispatcherOrchestrator{ + mc: mc, + dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), + } + processed := make(chan struct{}, 1) + for i := range orchestrator.shards { + orchestrator.shards[i] = newOrchestratorShard(func(msg *messaging.TargetMessage) { + processed <- struct{}{} + }) + orchestrator.shards[i].Run() + } + orchestrator.LocalFence() + for _, shard := range orchestrator.shards { + shard.Wait() + } + + cfID := common.NewChangeFeedIDWithName("cf", "default") + msg := messaging.NewSingleTargetMessage( + node.ID("to"), + messaging.DispatcherManagerManagerTopic, + &heartbeatpb.MaintainerCloseRequest{ChangefeedID: cfID.ToPB()}, + ) + require.NoError(t, orchestrator.RecvMaintainerRequest(context.Background(), msg)) + + select { + case <-processed: + require.FailNow(t, "local fence should drop new maintainer messages") + case <-time.After(50 * time.Millisecond): + } +} + +func TestDispatcherOrchestratorLocalFenceFencesManagersImmediately(t *testing.T) { + mc, _, stop := messaging.NewMessageCenterForTest(t) + defer stop() + + cfID := common.NewChangeFeedIDWithName("cf", "default") + manager := &dispatchermanager.DispatcherManager{} + orchestrator := &DispatcherOrchestrator{ + mc: mc, + dispatcherManagers: map[common.ChangeFeedID]*dispatchermanager.DispatcherManager{ + cfID: manager, + }, + shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), + } + for i := range orchestrator.shards { + orchestrator.shards[i] = newOrchestratorShard(func(msg *messaging.TargetMessage) {}) + orchestrator.shards[i].Run() + } + orchestrator.LocalFence() + for _, shard := range orchestrator.shards { + shard.Wait() + } + + err := manager.InitalizeTableTriggerEventDispatcher(nil) + require.True(t, dispatchermanager.IsWritePathClosedError(err)) +} + +func TestDispatcherOrchestratorLocalFenceFencesInitializingManagersImmediately(t *testing.T) { + mc, _, stop := messaging.NewMessageCenterForTest(t) + defer stop() + + cfID := common.NewChangeFeedIDWithName("cf", "default") + manager := &dispatchermanager.DispatcherManager{} + orchestrator := &DispatcherOrchestrator{ + mc: mc, + initializingDispatcherManagers: map[common.ChangeFeedID]*dispatchermanager.DispatcherManager{ + cfID: manager, + }, + dispatcherManagers: make(map[common.ChangeFeedID]*dispatchermanager.DispatcherManager), + shards: make([]*orchestratorShard, dispatcherOrchestratorShardCount), + } + for i := range orchestrator.shards { + orchestrator.shards[i] = newOrchestratorShard(func(msg *messaging.TargetMessage) {}) + orchestrator.shards[i].Run() + } + orchestrator.LocalFence() + for _, shard := range orchestrator.shards { + shard.Wait() + } + + err := manager.InitalizeTableTriggerEventDispatcher(nil) + require.True(t, dispatchermanager.IsWritePathClosedError(err)) +} diff --git a/pkg/common/context/app_context.go b/pkg/common/context/app_context.go index ead7d089b2..38f020b2fc 100644 --- a/pkg/common/context/app_context.go +++ b/pkg/common/context/app_context.go @@ -63,3 +63,13 @@ func GetService[T any](name string) T { v, _ := GetGlobalContext().serviceMap.Load(name) return v.(T) } + +func TryGetService[T any](name string) (T, bool) { + v, ok := GetGlobalContext().serviceMap.Load(name) + if !ok { + var zero T + return zero, false + } + service, ok := v.(T) + return service, ok +} diff --git a/pkg/errors/error.go b/pkg/errors/error.go index 4cf346f573..e15d05e818 100644 --- a/pkg/errors/error.go +++ b/pkg/errors/error.go @@ -778,6 +778,10 @@ var ( "failed to init table trigger dispatcher", errors.RFCCodeText("CDC:ErrChangefeedInitTableTriggerDispatcherFailed"), ) + ErrDispatcherManagerWritePathClosed = errors.Normalize( + "dispatcher manager write path is closed", + errors.RFCCodeText("CDC:ErrDispatcherManagerWritePathClosed"), + ) ErrDDLEventError = errors.Normalize( "ddl event meets error", errors.RFCCodeText("CDC:ErrDDLEventError"), diff --git a/pkg/sink/mysql/mysql_writer_dml.go b/pkg/sink/mysql/mysql_writer_dml.go index baa88c78ec..b5357e0b39 100644 --- a/pkg/sink/mysql/mysql_writer_dml.go +++ b/pkg/sink/mysql/mysql_writer_dml.go @@ -705,7 +705,10 @@ func (w *Writer) execDMLWithMaxRetries(dmls *preparedDMLs) error { failpoint.Return(err) }) - failpoint.Inject("MySQLSinkHangLongTime", func() { _ = util.Hang(w.ctx, time.Hour) }) + failpoint.Inject("MySQLSinkHangLongTime", func() { + log.Warn("inject MySQLSinkHangLongTime") + _ = util.Hang(w.ctx, time.Hour) + }) failpoint.Inject("MySQLDuplicateEntryError", func() { log.Warn("inject MySQLDuplicateEntryError") diff --git a/server/server.go b/server/server.go index 421f4f4762..7f0aee9b3b 100644 --- a/server/server.go +++ b/server/server.go @@ -46,6 +46,7 @@ import ( "github.com/pingcap/ticdc/pkg/upstream" "github.com/pingcap/ticdc/server/watcher" pd "github.com/tikv/pd/client" + clientv3 "go.etcd.io/etcd/client/v3" "go.etcd.io/etcd/client/v3/concurrency" "go.uber.org/zap" "golang.org/x/sync/errgroup" @@ -55,10 +56,15 @@ const ( closeServiceTimeout = 15 * time.Second cleanMetaDuration = 10 * time.Second oldArchCheckInterval = 100 * time.Millisecond + sessionWatchInterval = time.Second // GracefulShutdownTimeout is used to prevent the CDC process from hanging for an extended period due to certain modules don't exit immediately. GracefulShutdownTimeout = 30 * time.Second ) +type localFencer interface { + LocalFence() +} + // server represents the main TiCDC server with carefully orchestrated module lifecycle management. // // Module Startup Order (dependencies flow from top to bottom): @@ -135,6 +141,8 @@ type server struct { subModules []common.SubModule closed atomic.Bool + + localFenceOnce atomic.Bool } // New returns a new Server instance @@ -309,7 +317,9 @@ func (c *server) Run(ctx context.Context) error { }(sub) } - g, gctx := errgroup.WithContext(egctx) + groupCtx, cancelGroup := context.WithCancel(egctx) + defer cancelGroup() + g, gctx := errgroup.WithContext(groupCtx) // start all subCommonModules for _, sub := range c.nodeModules { func(m common.SubModule) { @@ -343,6 +353,14 @@ func (c *server) Run(ctx context.Context) error { return errors.Trace(err) } + fatalErrCh := make(chan error, 1) + go func() { + if err := c.runSessionWatchdog(gctx, sessionWatchInterval); err != nil { + fatalErrCh <- err + cancelGroup() + } + }() + // if it takes too long for all sub modules to exit, then exit directly to avoid hanging. ch := make(chan error, 1) go func() { @@ -353,10 +371,74 @@ func (c *server) Run(ctx context.Context) error { go func() { ch <- g.Wait() }() - err = <-ch + select { + case err = <-fatalErrCh: + case err = <-ch: + } return err } +func (c *server) runSessionWatchdog(ctx context.Context, checkInterval time.Duration) error { + if c.session == nil { + return nil + } + return c.watchEtcdSession(ctx, c.session.Done(), c.session.Lease(), checkInterval) +} + +func (c *server) watchEtcdSession( + ctx context.Context, + sessionDone <-chan struct{}, + leaseID clientv3.LeaseID, + checkInterval time.Duration, +) error { + ticker := time.NewTicker(checkInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return nil + case <-sessionDone: + if ctx.Err() != nil { + return nil + } + c.localFence("etcd session done") + return errors.ErrCaptureSuicide.GenWithStackByArgs() + case <-ticker.C: + ttl, err := c.EtcdClient.GetEtcdClient().TimeToLive(ctx, leaseID) + if err != nil { + if ctx.Err() != nil { + return nil + } + log.Warn("check etcd session ttl failed", zap.Error(err)) + continue + } + if ttl != nil && ttl.TTL == -1 { + c.localFence("etcd lease expired") + return errors.ErrCaptureSuicide.GenWithStackByArgs() + } + } + } +} + +func (c *server) localFence(reason string) { + if !c.localFenceOnce.CompareAndSwap(false, true) { + return + } + + log.Warn("local fence triggered", zap.String("reason", reason)) + + c.liveness.Store(liveness.CaptureDraining) + c.liveness.Store(liveness.CaptureStopping) + + fencer, ok := appctx.TryGetService[localFencer](appctx.DispatcherOrchestrator) + if !ok { + log.Warn("dispatcher orchestrator is not available when local fence triggered") + return + } + fencer.LocalFence() +} + // validCheck checks whether the environment is valid to start the server // return only when all the old-arch cdc capture is not running // old-arch cdc capture will return when receive the unknown etcd key diff --git a/server/server_session_watchdog_test.go b/server/server_session_watchdog_test.go new file mode 100644 index 0000000000..fea9788382 --- /dev/null +++ b/server/server_session_watchdog_test.go @@ -0,0 +1,88 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// See the License for the specific language governing permissions and +// limitations under the License. + +package server + +import ( + "context" + "sync/atomic" + "testing" + "time" + + "github.com/golang/mock/gomock" + appctx "github.com/pingcap/ticdc/pkg/common/context" + "github.com/pingcap/ticdc/pkg/errors" + "github.com/pingcap/ticdc/pkg/etcd" + "github.com/pingcap/ticdc/pkg/liveness" + "github.com/stretchr/testify/require" + clientv3 "go.etcd.io/etcd/client/v3" +) + +type testLocalFencer struct { + count atomic.Int32 +} + +func (f *testLocalFencer) LocalFence() { + f.count.Add(1) +} + +func TestSessionWatchdogFencesOnSessionDone(t *testing.T) { + fencer := &testLocalFencer{} + appctx.SetService(appctx.DispatcherOrchestrator, fencer) + + c := &server{} + sessionDone := make(chan struct{}) + close(sessionDone) + + err := c.watchEtcdSession(context.Background(), sessionDone, 1, time.Hour) + + require.True(t, errors.ErrCaptureSuicide.Equal(err), err) + require.Equal(t, int32(1), fencer.count.Load()) + require.Equal(t, liveness.CaptureStopping, c.liveness.Load()) +} + +func TestSessionWatchdogFencesOnExpiredLease(t *testing.T) { + fencer := &testLocalFencer{} + appctx.SetService(appctx.DispatcherOrchestrator, fencer) + + ctrl := gomock.NewController(t) + cdcEtcdClient := etcd.NewMockCDCEtcdClient(ctrl) + rawEtcdClient := etcd.NewMockClient(ctrl) + cdcEtcdClient.EXPECT().GetEtcdClient().Return(rawEtcdClient).AnyTimes() + rawEtcdClient.EXPECT(). + TimeToLive(gomock.Any(), clientv3.LeaseID(100)). + Return(&clientv3.LeaseTimeToLiveResponse{TTL: -1}, nil) + + c := &server{EtcdClient: cdcEtcdClient} + sessionDone := make(chan struct{}) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + err := c.watchEtcdSession(ctx, sessionDone, 100, time.Millisecond) + + require.True(t, errors.ErrCaptureSuicide.Equal(err), err) + require.Equal(t, int32(1), fencer.count.Load()) + require.Equal(t, liveness.CaptureStopping, c.liveness.Load()) +} + +func TestLocalFenceIsIdempotent(t *testing.T) { + fencer := &testLocalFencer{} + appctx.SetService(appctx.DispatcherOrchestrator, fencer) + + c := &server{} + c.localFence("first") + c.localFence("second") + + require.Equal(t, int32(1), fencer.count.Load()) + require.Equal(t, liveness.CaptureStopping, c.liveness.Load()) +} diff --git a/tests/integration_tests/capture_local_fence_on_session_done/conf/diff_config.toml b/tests/integration_tests/capture_local_fence_on_session_done/conf/diff_config.toml new file mode 100644 index 0000000000..543b483a11 --- /dev/null +++ b/tests/integration_tests/capture_local_fence_on_session_done/conf/diff_config.toml @@ -0,0 +1,29 @@ +# diff Configuration. + +check-thread-count = 4 + +export-fix-sql = true + +check-struct-only = false + +[task] + output-dir = "/tmp/tidb_cdc_test/capture_local_fence_on_session_done/sync_diff/output" + + source-instances = ["mysql1"] + + target-instance = "tidb0" + + target-check-tables = ["capture_local_fence_on_session_done.?*"] + +[data-sources] +[data-sources.mysql1] + host = "127.0.0.1" + port = 4000 + user = "root" + password = "" + +[data-sources.tidb0] + host = "127.0.0.1" + port = 3306 + user = "root" + password = "" diff --git a/tests/integration_tests/capture_local_fence_on_session_done/run.sh b/tests/integration_tests/capture_local_fence_on_session_done/run.sh new file mode 100755 index 0000000000..79d105b208 --- /dev/null +++ b/tests/integration_tests/capture_local_fence_on_session_done/run.sh @@ -0,0 +1,170 @@ +#!/bin/bash + +set -eu + +CUR=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +source $CUR/../_utils/test_prepare +WORK_DIR=$OUT_DIR/$TEST_NAME +CDC_BINARY=cdc.test +SINK_TYPE=$1 + +DB_NAME="capture_local_fence_on_session_done" +TABLE_NAME="t" +CAPTURE1_PORT=${CDC_PORT:-8300} +CAPTURE2_PORT=${CAPTURE2_PORT:-8301} +CAPTURE1_LISTEN_ADDR="127.0.0.1:${CAPTURE1_PORT}" +CAPTURE2_LISTEN_ADDR="127.0.0.1:${CAPTURE2_PORT}" +export CDC_PORT=$CAPTURE1_PORT +HANG_FAILPOINT="github.com/pingcap/ticdc/pkg/sink/mysql/MySQLSinkHangLongTime" + +function get_capture_addr_by_port() { + local api_addr=$1 + local port=$2 + local suffix=":$port" + + for ((i = 0; i < 30; i++)); do + local addr + addr=$(curl -sS --connect-timeout 2 --max-time 5 "http://${api_addr}/api/v2/captures" | + jq -r --arg suffix "$suffix" '.items[] | select(.address | endswith($suffix)) | .address' | head -n1) + if [ -n "$addr" ] && [ "$addr" != "null" ]; then + echo "$addr" + return 0 + fi + sleep 2 + done + + echo "capture with port $port not found" >&2 + return 1 +} + +function get_capture_id_by_addr() { + local api_addr=$1 + local target_addr=$2 + curl -sS --connect-timeout 2 --max-time 5 "http://${api_addr}/api/v2/captures" | + jq -r --arg addr "$target_addr" '.items[] | select(.address==$addr) | .id' | head -n1 +} + +function get_table_node_id() { + local api_addr=$1 + local changefeed_id=$2 + local table_id=$3 + curl -sS --connect-timeout 2 --max-time 5 "http://${api_addr}/api/v2/changefeeds/${changefeed_id}/tables?keyspace=$KEYSPACE_NAME" | + jq -r --argjson tid "$table_id" '.items[] | select(.table_ids | index($tid)) | .node_id' | head -n1 +} + +function wait_for_table_on_addr() { + local api_addr=$1 + local changefeed_id=$2 + local table_id=$3 + local target_addr=$4 + + for ((i = 0; i < 30; i++)); do + local target_id + target_id=$(get_capture_id_by_addr "$api_addr" "$target_addr") + if [ -z "$target_id" ] || [ "$target_id" == "null" ]; then + sleep 2 + continue + fi + + local node_id + node_id=$(get_table_node_id "$api_addr" "$changefeed_id" "$table_id") + if [ "$node_id" == "$target_id" ]; then + return 0 + fi + sleep 2 + done + + echo "table $table_id not scheduled on $target_addr" >&2 + return 1 +} + +function revoke_capture_lease() { + local capture_id=$1 + local capture_key="/tidb/cdc/default/__cdc_meta__/capture/${capture_id}" + local lease + lease=$(ETCDCTL_API=3 etcdctl get "$capture_key" -w json | grep -o 'lease":[0-9]*' | awk -F: '{print $2}') + if [ -z "$lease" ]; then + echo "failed to get etcd lease for capture $capture_id" >&2 + return 1 + fi + + local lease_hex + lease_hex=$(printf '%x\n' "$lease") + ETCDCTL_API=3 etcdctl lease revoke "$lease_hex" +} + +function wait_for_downstream_row_count() { + local expected=$1 + for ((i = 0; i < 30; i++)); do + local count + count=$(mysql -h${DOWN_TIDB_HOST} -P${DOWN_TIDB_PORT} -uroot -N \ + -e "select count(*) from ${DB_NAME}.${TABLE_NAME};" 2>/dev/null || true) + if [ "$count" == "$expected" ]; then + return 0 + fi + sleep 2 + done + + echo "downstream row count is not $expected" >&2 + return 1 +} + +function run() { + # This case depends on the MySQL sink DML failpoint. + if [ "$SINK_TYPE" != "mysql" ]; then + return + fi + + rm -rf $WORK_DIR && mkdir -p $WORK_DIR + start_tidb_cluster --workdir $WORK_DIR + + pd_addr="http://$UP_PD_HOST_1:$UP_PD_PORT_1" + run_cdc_server --workdir $WORK_DIR --binary $CDC_BINARY --pd $pd_addr --logsuffix 1 --addr "$CAPTURE1_LISTEN_ADDR" + export GO_FAILPOINTS="${HANG_FAILPOINT}=return(true)" + run_cdc_server --workdir $WORK_DIR --binary $CDC_BINARY --pd $pd_addr --logsuffix 2 --addr "$CAPTURE2_LISTEN_ADDR" + export GO_FAILPOINTS='' + capture2_addr=$(get_capture_addr_by_port "$CAPTURE1_LISTEN_ADDR" "${CAPTURE2_LISTEN_ADDR#*:}") + + SINK_URI="mysql://normal:123456@127.0.0.1:3306/?max-txn-row=1&worker-count=1" + + run_sql "CREATE DATABASE ${DB_NAME};" ${UP_TIDB_HOST} ${UP_TIDB_PORT} + run_sql "CREATE TABLE ${DB_NAME}.${TABLE_NAME} (id int primary key, v int);" ${UP_TIDB_HOST} ${UP_TIDB_PORT} + run_sql "CREATE DATABASE ${DB_NAME};" ${DOWN_TIDB_HOST} ${DOWN_TIDB_PORT} + run_sql "CREATE TABLE ${DB_NAME}.${TABLE_NAME} (id int primary key, v int);" ${DOWN_TIDB_HOST} ${DOWN_TIDB_PORT} + start_ts=$(run_cdc_cli_tso_query ${UP_PD_HOST_1} ${UP_PD_PORT_1}) + + changefeed_id=$(cdc_cli_changefeed create --pd=$pd_addr --start-ts=$start_ts --sink-uri="$SINK_URI" | grep '^ID:' | head -n1 | awk '{print $2}') + table_id=$(get_table_id "$DB_NAME" "$TABLE_NAME") + + move_table_with_retry "$capture2_addr" $table_id "$changefeed_id" 10 + wait_for_table_on_addr "$CAPTURE1_LISTEN_ADDR" "$changefeed_id" "$table_id" "$capture2_addr" + + capture2_id=$(get_capture_id_by_addr "$CAPTURE1_LISTEN_ADDR" "$capture2_addr") + if [ -z "$capture2_id" ] || [ "$capture2_id" == "null" ]; then + echo "failed to get capture id for $capture2_addr" >&2 + exit 1 + fi + cdc2_pid=$(get_cdc_pid "${CAPTURE2_LISTEN_ADDR%:*}" "${CAPTURE2_LISTEN_ADDR#*:}") + + run_sql "INSERT INTO ${DB_NAME}.${TABLE_NAME} VALUES (1, 10), (2, 20), (3, 30), (4, 40);" ${UP_TIDB_HOST} ${UP_TIDB_PORT} + ensure 30 "grep -Eq 'inject MySQLSinkHangLongTime' '$WORK_DIR/cdc2.log'" + wait_for_downstream_row_count 0 + + revoke_capture_lease "$capture2_id" + + ensure 15 "! ps -p $cdc2_pid" + check_logs_contains $WORK_DIR "local fence triggered" 2 + check_logs_contains $WORK_DIR "dispatcher orchestrator local fence triggered" 2 + check_logs_contains $WORK_DIR "stopping dispatcher manager write path" 2 + check_logs_contains $WORK_DIR "server closed" 2 + + run_sql "INSERT INTO ${DB_NAME}.${TABLE_NAME} VALUES (5, 50), (6, 60);" ${UP_TIDB_HOST} ${UP_TIDB_PORT} + check_sync_diff $WORK_DIR $CUR/conf/diff_config.toml + + cleanup_process $CDC_BINARY +} + +trap 'stop_test $WORK_DIR' EXIT +run "$@" +check_logs $WORK_DIR +echo "[$(date)] <<<<<< run test case $TEST_NAME success! >>>>>>" diff --git a/tests/integration_tests/capture_session_done_during_task/run.sh b/tests/integration_tests/capture_session_done_during_task/run.sh index 643d58c906..71b5d0e5a6 100755 --- a/tests/integration_tests/capture_session_done_during_task/run.sh +++ b/tests/integration_tests/capture_session_done_during_task/run.sh @@ -56,7 +56,8 @@ function run() { # ensure server exit ensure 30 "! ps -p $cdc_pid" - check_logs_contains $WORK_DIR "the etcd session is done" + check_logs_contains $WORK_DIR "local fence triggered" + ensure 30 "grep -Eq 'etcd lease expired|etcd session done' '$WORK_DIR/cdc.log'" check_logs_contains $WORK_DIR "server closed" echo "cdc server already exit" @@ -75,6 +76,6 @@ function run() { } trap 'stop_test $WORK_DIR' EXIT -run $* +run "$@" check_logs $WORK_DIR echo "[$(date)] <<<<<< run test case $TEST_NAME success! >>>>>>" diff --git a/tests/integration_tests/run_light_it_in_ci.sh b/tests/integration_tests/run_light_it_in_ci.sh index fbb07d8d2c..3a3ec7488c 100755 --- a/tests/integration_tests/run_light_it_in_ci.sh +++ b/tests/integration_tests/run_light_it_in_ci.sh @@ -38,7 +38,7 @@ mysql_groups=( # G02 'new_ci_collation safe_mode savepoint fail_over_ddl_C unsplittable_tables' # G03 - 'capture_suicide_while_balance_table kv_client_stream_reconnect fail_over_ddl_D ddl_default_current_timestamp' + 'capture_suicide_while_balance_table capture_local_fence_on_session_done kv_client_stream_reconnect fail_over_ddl_D ddl_default_current_timestamp' # G04 'multi_capture ci_collation_compatibility resourcecontrol fail_over_ddl_E' # G05