@@ -125,7 +125,7 @@ func newRegionRequestWorker(
125125
126126func (s * regionRequestWorker ) Run (ctx context.Context ) error {
127127 handleStreamFailure := func (firstReq * regionReq , regionErr error ) {
128- // Stream failure recovery :
128+ // Stream failure handle cases :
129129 // - tracker: requests already sent to this stream.
130130 // - firstReq: popped from admission for this stream, but not necessarily
131131 // added to tracker yet if the stream fails before sendRegionRequest calls
@@ -217,8 +217,7 @@ func (s *regionRequestWorker) runStream(ctx context.Context, firstReq *regionReq
217217 zap .Error (err ))
218218 }()
219219
220- g , gctx := errgroup .WithContext (ctx )
221- conn , err := Connect (gctx , s .upstream .credential , s .storeAddr )
220+ conn , err := Connect (ctx , s .upstream .credential , s .storeAddr )
222221 if err != nil {
223222 log .Warn ("region request worker create grpc stream failed" ,
224223 zap .Uint64 ("workerID" , s .workerID ),
@@ -234,6 +233,7 @@ func (s *regionRequestWorker) runStream(ctx context.Context, firstReq *regionReq
234233 }
235234 defer func () { _ = conn .Conn .Close () }()
236235
236+ g , gctx := errgroup .WithContext (ctx )
237237 g .Go (func () error { return s .receiveAndDispatchChangeEvents (conn ) })
238238 g .Go (func () error { return s .processRegionSendTask (gctx , conn , firstReq ) })
239239
@@ -285,13 +285,13 @@ func (s *regionRequestWorker) receiveAndDispatchChangeEvents(conn *ConnAndClient
285285}
286286
287287func (s * regionRequestWorker ) dispatchRegionChangeEvents (events []* cdcpb.Event ) {
288- for _ , cdcEvent := range events {
289- regionID := cdcEvent .RegionId
290- subscriptionID := SubscriptionID (cdcEvent .RequestId )
288+ for _ , event := range events {
289+ regionID := event .RegionId
290+ subscriptionID := SubscriptionID (event .RequestId )
291291 state := s .tracker .Get (subscriptionID , regionID )
292292 if state != nil {
293293 regionEvent := regionEvent {states : []* regionFeedState {state }}
294- switch eventData := cdcEvent .Event .(type ) {
294+ switch eventData := event .Event .(type ) {
295295 case * cdcpb.Event_Entries_ :
296296 if eventData == nil {
297297 log .Warn ("region request worker receives a region event with nil entries, ignore it" ,
@@ -307,7 +307,7 @@ func (s *regionRequestWorker) dispatchRegionChangeEvents(events []*cdcpb.Event)
307307 log .Debug ("region request worker receives a region error" ,
308308 zap .Uint64 ("workerID" , s .workerID ),
309309 zap .Uint64 ("subscriptionID" , uint64 (subscriptionID )),
310- zap .Uint64 ("regionID" , cdcEvent .RegionId ),
310+ zap .Uint64 ("regionID" , event .RegionId ),
311311 zap .Any ("error" , eventData .Error ))
312312 state .markStopped (& eventError {err : eventData .Error })
313313 s .eventSink .Push (subscriptionID , regionEvent )
@@ -317,23 +317,23 @@ func (s *regionRequestWorker) dispatchRegionChangeEvents(events []*cdcpb.Event)
317317 case * cdcpb.Event_LongTxn_ :
318318 continue
319319 default :
320- log .Panic ("unknown event type" , zap .Any ("event" , cdcEvent ))
320+ log .Panic ("unknown event type" , zap .Any ("event" , event ))
321321 }
322322 s .eventSink .Push (subscriptionID , regionEvent )
323323 continue
324324 }
325325
326- switch cdcEvent .Event .(type ) {
326+ switch event .Event .(type ) {
327327 case * cdcpb.Event_Error :
328328 log .Debug ("region request worker receives an error for a stale region, ignore it" ,
329329 zap .Uint64 ("workerID" , s .workerID ),
330330 zap .Uint64 ("subscriptionID" , uint64 (subscriptionID )),
331- zap .Uint64 ("regionID" , cdcEvent .RegionId ))
331+ zap .Uint64 ("regionID" , event .RegionId ))
332332 default :
333333 log .Warn ("region request worker receives a region event for an untracked region" ,
334334 zap .Uint64 ("workerID" , s .workerID ),
335335 zap .Uint64 ("subscriptionID" , uint64 (subscriptionID )),
336- zap .Uint64 ("regionID" , cdcEvent .RegionId ))
336+ zap .Uint64 ("regionID" , event .RegionId ))
337337 }
338338 }
339339}
0 commit comments