summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com>2020-12-07 09:16:38 +0800
committerGitHub <noreply@github.com>2020-12-07 09:16:38 +0800
commit97b47f8db77fdb6ba4266ef647387795b6e7cb5d (patch)
tree783d9690621a5979ed335def6da3d9f255bd0a56
parentMerge pull request #2377 from daixiang0/automated-cherry-pick-of-#2334-upstre... (diff)
parentfix message send problem (diff)
downloadkubeedge-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-xcloud/pkg/cloudhub/handler/messagehandler.go18
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
}