Skip to content

Commit 7448aa1

Browse files
guyinyouguyinyou
andauthored
[Go] Fix data race in receiveMessage for push and simple consumers (#1275)
The err and resps variables were shared between the main goroutine and a spawned goroutine without synchronization, causing a data race detectable by Go's race detector. The goroutine wrote to both variables while the select handler read them concurrently. Move resps and err into the goroutine as local variables, and pass results back through a typed channel (receiveResult struct) instead. This also fixes push_consumer.go not signaling the done channel on non-EOF errors, which previously caused it to wait for context timeout instead of returning the actual error immediately. Affected: - golang/push_consumer.go: receiveMessage() - golang/simple_consumer.go: receiveMessage() Co-authored-by: guyinyou <guyinyou.gyy@alibaba-inc.com>
1 parent 882c5be commit 7448aa1

2 files changed

Lines changed: 34 additions & 30 deletions

File tree

golang/push_consumer.go

Lines changed: 17 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -206,7 +206,6 @@ func (pc *defaultPushConsumer) GetGroupName() string {
206206
}
207207

208208
func (pc *defaultPushConsumer) receiveMessage(ctx context.Context, request *v2.ReceiveMessageRequest, messageQueue *v2.MessageQueue, timeout time.Duration) ([]*MessageView, error) {
209-
var err error
210209
ctx = pc.cli.Sign(ctx)
211210
ctx, cancel := context.WithTimeout(ctx, timeout)
212211
defer cancel()
@@ -215,34 +214,38 @@ func (pc *defaultPushConsumer) receiveMessage(ctx context.Context, request *v2.R
215214
if err != nil {
216215
return nil, err
217216
}
218-
done := make(chan bool, 1)
217+
type receiveResult struct {
218+
responses []*v2.ReceiveMessageResponse
219+
err error
220+
}
221+
done := make(chan receiveResult, 1)
219222

220-
resps := make([]*v2.ReceiveMessageResponse, 0)
221223
go func() {
224+
var resps []*v2.ReceiveMessageResponse
225+
var recvErr error
222226
for {
223227
var resp *v2.ReceiveMessageResponse
224-
resp, err = receiveMessageClient.Recv()
225-
if err == io.EOF {
226-
done <- true
227-
defer close(done)
228+
resp, recvErr = receiveMessageClient.Recv()
229+
if recvErr == io.EOF {
230+
done <- receiveResult{responses: resps}
228231
break
229232
}
230-
if err != nil {
231-
pc.cli.log.Errorf("pushConsumer recv msg err=%v, requestId=%s", err, utils.GetRequestID(ctx))
233+
if recvErr != nil {
234+
pc.cli.log.Errorf("pushConsumer recv msg err=%v, requestId=%s", recvErr, utils.GetRequestID(ctx))
235+
done <- receiveResult{err: recvErr}
232236
break
233237
}
234238
sugarBaseLogger.Debugf("receiveMessage response: %v", resp)
235239
resps = append(resps, resp)
236240
}
237-
cancel()
238241
}()
239242
select {
240243
case <-ctx.Done():
241244
// timeout
242245
return nil, fmt.Errorf("[error] CODE=DEADLINE_EXCEEDED")
243-
case <-done:
244-
if err != nil && err != io.EOF {
245-
return nil, err
246+
case result := <-done:
247+
if result.err != nil {
248+
return nil, result.err
246249
}
247250
messageViewList := make([]*MessageView, 0)
248251
status := &v2.Status{
@@ -251,7 +254,7 @@ func (pc *defaultPushConsumer) receiveMessage(ctx context.Context, request *v2.R
251254
}
252255
var deliveryTimestamp *timestamppb.Timestamp
253256
messageList := make([]*v2.Message, 0)
254-
for _, resp := range resps {
257+
for _, resp := range result.responses {
255258
switch r := resp.GetContent().(type) {
256259
case *v2.ReceiveMessageResponse_Status:
257260
status = r.Status

golang/simple_consumer.go

Lines changed: 17 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -210,7 +210,6 @@ func (sc *defaultSimpleConsumer) GetGroupName() string {
210210
}
211211

212212
func (sc *defaultSimpleConsumer) receiveMessage(ctx context.Context, request *v2.ReceiveMessageRequest, messageQueue *v2.MessageQueue, timeout time.Duration) ([]*MessageView, error) {
213-
var err error
214213
ctx = sc.cli.Sign(ctx)
215214
ctx, cancel := context.WithTimeout(ctx, timeout)
216215
defer cancel()
@@ -219,36 +218,38 @@ func (sc *defaultSimpleConsumer) receiveMessage(ctx context.Context, request *v2
219218
if err != nil {
220219
return nil, err
221220
}
222-
done := make(chan bool, 1)
221+
type receiveResult struct {
222+
responses []*v2.ReceiveMessageResponse
223+
err error
224+
}
225+
done := make(chan receiveResult, 1)
223226

224-
resps := make([]*v2.ReceiveMessageResponse, 0)
225227
go func() {
228+
var resps []*v2.ReceiveMessageResponse
229+
var recvErr error
226230
for {
227231
var resp *v2.ReceiveMessageResponse
228-
resp, err = receiveMessageClient.Recv()
229-
if err == io.EOF {
230-
done <- true
231-
defer close(done)
232+
resp, recvErr = receiveMessageClient.Recv()
233+
if recvErr == io.EOF {
234+
done <- receiveResult{responses: resps}
232235
break
233236
}
234-
if err != nil {
235-
sc.cli.log.Errorf("simpleConsumer recv msg err=%v, requestId=%s", err, utils.GetRequestID(ctx))
236-
done <- true
237-
defer close(done)
237+
if recvErr != nil {
238+
sc.cli.log.Errorf("simpleConsumer recv msg err=%v, requestId=%s", recvErr, utils.GetRequestID(ctx))
239+
done <- receiveResult{err: recvErr}
238240
break
239241
}
240242
sugarBaseLogger.Debugf("receiveMessage response: %v", resp)
241243
resps = append(resps, resp)
242244
}
243-
cancel()
244245
}()
245246
select {
246247
case <-ctx.Done():
247248
// timeout
248249
return nil, fmt.Errorf("[error] CODE=DEADLINE_EXCEEDED")
249-
case <-done:
250-
if err != nil && err != io.EOF {
251-
return nil, err
250+
case result := <-done:
251+
if result.err != nil {
252+
return nil, result.err
252253
}
253254
messageViewList := make([]*MessageView, 0)
254255
status := &v2.Status{
@@ -257,7 +258,7 @@ func (sc *defaultSimpleConsumer) receiveMessage(ctx context.Context, request *v2
257258
}
258259
var deliveryTimestamp *timestamppb.Timestamp
259260
messageList := make([]*v2.Message, 0)
260-
for _, resp := range resps {
261+
for _, resp := range result.responses {
261262
switch r := resp.GetContent().(type) {
262263
case *v2.ReceiveMessageResponse_Status:
263264
status = r.Status

0 commit comments

Comments
 (0)