fix(msggateway): avoid closed channel panic in super group push

pull/3767/head
buvidk 2 months ago
parent de0f1eec8a
commit badf05f6b1

@ -175,22 +175,22 @@ func (s *Server) SuperGroupOnlineBatchPushOneMsg(ctx context.Context, req *msgga
ch := make(chan *msggateway.SingleMsgToUserResults, len(req.PushToUserIDs)) ch := make(chan *msggateway.SingleMsgToUserResults, len(req.PushToUserIDs))
var count atomic.Int64 var count atomic.Int64
count.Add(int64(len(req.PushToUserIDs))) count.Add(int64(len(req.PushToUserIDs)))
for i := range req.PushToUserIDs { pushResult := func(result *msggateway.SingleMsgToUserResults) {
userID := req.PushToUserIDs[i] ch <- result
err := s.queue.PushCtx(ctx, func() {
ch <- s.pushToUser(ctx, userID, req.MsgData)
if count.Add(-1) == 0 { if count.Add(-1) == 0 {
close(ch) close(ch)
} }
}
for i := range req.PushToUserIDs {
userID := req.PushToUserIDs[i]
err := s.queue.PushCtx(ctx, func() {
pushResult(s.pushToUser(ctx, userID, req.MsgData))
}) })
if err != nil { if err != nil {
if count.Add(-1) == 0 {
close(ch)
}
log.ZError(ctx, "pushToUser MemoryQueue failed", err, "userID", userID) log.ZError(ctx, "pushToUser MemoryQueue failed", err, "userID", userID)
ch <- &msggateway.SingleMsgToUserResults{ pushResult(&msggateway.SingleMsgToUserResults{
UserID: userID, UserID: userID,
} })
} }
} }
resp := &msggateway.OnlineBatchPushOneMsgResp{ resp := &msggateway.OnlineBatchPushOneMsgResp{

Loading…
Cancel
Save