diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2020-12-07 09:16:38 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2020-12-07 09:16:38 +0800 |
| commit | 97b47f8db77fdb6ba4266ef647387795b6e7cb5d (patch) | |
| tree | 783d9690621a5979ed335def6da3d9f255bd0a56 | |
| parent | Merge pull request #2377 from daixiang0/automated-cherry-pick-of-#2334-upstre... (diff) | |
| parent | fix message send problem (diff) | |
| download | kubeedge-origin/release-1.5.tar.gz | |
Merge pull request #2399 from threestoneliu/automated-cherry-pick-of-#2392-upstream-release-1.5origin/release-1.5
Automated cherry pick of #2392: fix message send problem
| -rwxr-xr-x | cloud/pkg/cloudhub/handler/messagehandler.go | 18 |
1 files changed, 10 insertions, 8 deletions
diff --git a/cloud/pkg/cloudhub/handler/messagehandler.go b/cloud/pkg/cloudhub/handler/messagehandler.go index f4746458d..b8edac930 100755 --- a/cloud/pkg/cloudhub/handler/messagehandler.go +++ b/cloud/pkg/cloudhub/handler/messagehandler.go @@ -425,18 +425,20 @@ func (mh *MessageHandle) MessageWriteLoop(info *model.HubInfo, stopServe chan Ex obj, exist, _ := nodeStore.GetByKey(key.(string)) if !exist { klog.Errorf("nodeStore for node %s doesn't exist", info.NodeID) + nodeQueue.Done(key) continue } msg := obj.(*beehiveModel.Message) if !model.IsToEdge(msg) { klog.Infof("skip only to cloud event for node %s, %s, content %s", info.NodeID, dumpMessageMetadata(msg), msg.Content) + nodeQueue.Done(key) continue } klog.V(4).Infof("event to send for node %s, %s, content %s", info.NodeID, dumpMessageMetadata(msg), msg.Content) copyMsg := deepcopy(msg) - trimMessage(msg) + trimMessage(copyMsg) for { conn, ok := mh.nodeConns.Load(info.NodeID) @@ -444,10 +446,10 @@ func (mh *MessageHandle) MessageWriteLoop(info *model.HubInfo, stopServe chan Ex time.Sleep(time.Second * 2) continue } - err := mh.sendMsg(conn.(hubio.CloudHubIO), info, msg, copyMsg, nodeStore) + err := mh.sendMsg(conn.(hubio.CloudHubIO), info, copyMsg, msg, nodeStore) if err != nil { klog.Errorf("Failed to send event to node: %s, affected event: %s, err: %s", - info.NodeID, dumpMessageMetadata(msg), err.Error()) + info.NodeID, dumpMessageMetadata(copyMsg), err.Error()) nodeQueue.AddRateLimited(key.(string)) time.Sleep(time.Second * 2) } @@ -459,9 +461,9 @@ func (mh *MessageHandle) MessageWriteLoop(info *model.HubInfo, stopServe chan Ex } } -func (mh *MessageHandle) sendMsg(hi hubio.CloudHubIO, info *model.HubInfo, msg, copyMsg *beehiveModel.Message, nodeStore cache.Store) error { +func (mh *MessageHandle) sendMsg(hi hubio.CloudHubIO, info *model.HubInfo, copyMsg, msg *beehiveModel.Message, nodeStore cache.Store) error { ackChan := make(chan struct{}) - mh.MessageAcks.Store(msg.GetID(), ackChan) + mh.MessageAcks.Store(copyMsg.GetID(), ackChan) // initialize timer and retry count for sending message var ( @@ -469,7 +471,7 @@ func (mh *MessageHandle) sendMsg(hi hubio.CloudHubIO, info *model.HubInfo, msg, retryInterval time.Duration = 5 ) ticker := time.NewTimer(retryInterval * time.Second) - err := mh.send(hi, info, msg) + err := mh.send(hi, info, copyMsg) if err != nil { return err } @@ -478,13 +480,13 @@ LOOP: for { select { case <-ackChan: - mh.saveSuccessPoint(copyMsg, info, nodeStore) + mh.saveSuccessPoint(msg, info, nodeStore) break LOOP case <-ticker.C: if retry == 4 { break LOOP } - err := mh.send(hi, info, msg) + err := mh.send(hi, info, copyMsg) if err != nil { return err } |
