From 306c5fdab6f8df37ee133e58cbde857c979c7088 Mon Sep 17 00:00:00 2001 From: lidezhu Date: Wed, 5 Aug 2026 21:49:23 +0800 Subject: [PATCH] remove use buffer --- logservice/logpuller/region_event_sink.go | 3 -- .../logpuller/region_failure_handler.go | 30 +++++++++++++++++-- .../region_request_scheduler_test.go | 5 ++++ .../logpuller/subscription_client_test.go | 20 ++++++++++++- 4 files changed, 51 insertions(+), 7 deletions(-) diff --git a/logservice/logpuller/region_event_sink.go b/logservice/logpuller/region_event_sink.go index 64c837f320..e79a57ebbe 100644 --- a/logservice/logpuller/region_event_sink.go +++ b/logservice/logpuller/region_event_sink.go @@ -39,9 +39,6 @@ func newRegionEventSink( option := dynstream.NewOption() // Note: it is max batch size of the kv sent from tikv(not committed rows) option.BatchCount = 1024 - // TODO: Set `UseBuffer` to true until we refactor the `regionEventHandler.Handle` method so that it doesn't call any method of the dynamic stream. Currently, if `UseBuffer` is set to false, there will be a deadlock: - // ds.handleLoop fetch events from `ch` -> regionEventHandler.Handle -> ds.RemovePath -> send event to `ch` - option.UseBuffer = true ds := dynstream.NewParallelDynamicStream( "log-puller", ®ionEventHandler{eventSink: sink, failureHandler: failureHandler}, diff --git a/logservice/logpuller/region_failure_handler.go b/logservice/logpuller/region_failure_handler.go index 64568df09c..bb35dc2c8d 100644 --- a/logservice/logpuller/region_failure_handler.go +++ b/logservice/logpuller/region_failure_handler.go @@ -194,7 +194,8 @@ func (r *regionFailureHandler) Report(errInfo regionErrorInfo) { if errInfo.subscribedSpan.rangeLock.UnlockRange( errInfo.span.StartKey, errInfo.span.EndKey, errInfo.verID.GetID(), errInfo.verID.GetVer(), errInfo.resolvedTs()) { - r.onTableDrained(errInfo.subscribedSpan) + // Defer span cleanup to Run so Report never calls back into dynstream. + r.cache.addDrainedSpan(errInfo.subscribedSpan) return } r.cache.add(errInfo) @@ -206,6 +207,9 @@ func (r *regionFailureHandler) Run(ctx context.Context) error { defer r.cancelRecoveries() handleCachedErrors := func() error { + for _, span := range r.cache.popDrainedSpans() { + r.onTableDrained(span) + } for { batch := r.cache.popBatch(errCacheBatchSize) for _, errInfo := range batch { @@ -383,8 +387,9 @@ func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionEr type errCache struct { sync.Mutex - cache []regionErrorInfo - notify chan struct{} + cache []regionErrorInfo + drainedSpans []*subscribedSpan + notify chan struct{} } const errCacheBatchSize = 1024 @@ -400,12 +405,31 @@ func (e *errCache) add(errInfo regionErrorInfo) { e.Lock() defer e.Unlock() e.cache = append(e.cache, errInfo) + e.signal() +} + +func (e *errCache) addDrainedSpan(span *subscribedSpan) { + e.Lock() + defer e.Unlock() + e.drainedSpans = append(e.drainedSpans, span) + e.signal() +} + +func (e *errCache) signal() { select { case e.notify <- struct{}{}: default: } } +func (e *errCache) popDrainedSpans() []*subscribedSpan { + e.Lock() + defer e.Unlock() + drainedSpans := e.drainedSpans + e.drainedSpans = nil + return drainedSpans +} + func (e *errCache) popBatch(limit int) []regionErrorInfo { e.Lock() defer e.Unlock() diff --git a/logservice/logpuller/region_request_scheduler_test.go b/logservice/logpuller/region_request_scheduler_test.go index 6180d362df..0073f88636 100644 --- a/logservice/logpuller/region_request_scheduler_test.go +++ b/logservice/logpuller/region_request_scheduler_test.go @@ -207,6 +207,10 @@ func TestRegionRequestSchedulerSkipsStoppedSubscriptionBeforeCreatingStore(t *te handler := newRegionFailureHandler(nil, func(rt *subscribedSpan) { drainedCh <- rt }, nil, nil) + handlerErrCh := make(chan error, 1) + go func() { + handlerErrCh <- handler.Run(ctx) + }() scheduler := ®ionRequestScheduler{ upstream: &upstreamHandle{ pd: pdClient, @@ -244,4 +248,5 @@ func TestRegionRequestSchedulerSkipsStoppedSubscriptionBeforeCreatingStore(t *te cancel() require.ErrorIs(t, <-errCh, context.Canceled) + require.ErrorIs(t, <-handlerErrCh, context.Canceled) } diff --git a/logservice/logpuller/subscription_client_test.go b/logservice/logpuller/subscription_client_test.go index fa42a096c1..b7706954e6 100644 --- a/logservice/logpuller/subscription_client_test.go +++ b/logservice/logpuller/subscription_client_test.go @@ -393,7 +393,25 @@ func TestRegionFailureHandlerQueuesCanceledError(t *testing.T) { }, &requestCancelledErr{})) require.Len(t, client.failureHandler.cache.cache, 1) - require.Nil(t, client.spanRegistry.Get(span.subID)) + require.Same(t, span, client.spanRegistry.Get(span.subID)) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + runDone := make(chan error, 1) + go func() { + runDone <- client.failureHandler.Run(ctx) + }() + + require.Eventually(t, func() bool { + return client.spanRegistry.Get(span.subID) == nil + }, time.Second, 10*time.Millisecond) + cancel() + select { + case err := <-runDone: + require.ErrorIs(t, err, context.Canceled) + case <-time.After(time.Second): + t.Fatal("failure handler did not exit after context cancellation") + } } type mockDynamicStream struct{}