diff options
| -rw-r--r-- | cloud/pkg/cloudhub/dispatcher/message_dispatcher.go | 5 | ||||
| -rw-r--r-- | cloud/pkg/cloudhub/dispatcher/message_dispatcher_test.go | 14 | ||||
| -rw-r--r-- | cloud/pkg/edgecontroller/controller/upstream.go | 49 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/configmap.go | 5 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/secret.go | 5 | ||||
| -rw-r--r-- | edge/pkg/metamanager/process.go | 21 | ||||
| -rw-r--r-- | edge/pkg/metamanager/process_test.go | 22 |
7 files changed, 85 insertions, 36 deletions
diff --git a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go index 364fc3c83..622f5f80a 100644 --- a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go +++ b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go @@ -375,6 +375,11 @@ func noAckRequired(msg *beehivemodel.Message) bool { if ok && content == commonconst.MessageSuccessfulContent { return true } + // `error message` is not required to ack + _, ok = msg.Content.(error) + if ok { + return true + } fallthrough default: if msg.GetSource() == modules.EdgeControllerModuleName { diff --git a/cloud/pkg/cloudhub/dispatcher/message_dispatcher_test.go b/cloud/pkg/cloudhub/dispatcher/message_dispatcher_test.go index 0dfa14524..ab26f8a9f 100644 --- a/cloud/pkg/cloudhub/dispatcher/message_dispatcher_test.go +++ b/cloud/pkg/cloudhub/dispatcher/message_dispatcher_test.go @@ -17,6 +17,7 @@ limitations under the License. package dispatcher import ( + "fmt" "reflect" "testing" @@ -72,11 +73,6 @@ func TestNoAckRequired(t *testing.T) { want: true, }, { - name: "applicationResponse message", - message: beehivemodel.NewMessage("").SetResourceOperation("/node/edge-test/ignore/Application/ignore", "applicationResponse"), - want: true, - }, - { name: "user data message", message: beehivemodel.NewMessage("router").SetRoute("", "user"), want: true, @@ -87,13 +83,13 @@ func TestNoAckRequired(t *testing.T) { want: true, }, { - name: "response ok message", - message: beehivemodel.NewMessage("").SetResourceOperation("node/edge-node/default/node/edge-node", "response").FillBody("OK"), + name: "node message", + message: beehivemodel.NewMessage("").SetResourceOperation("node/edge-node/default/node/edge-node", "response").SetRoute("edgecontroller", "resource"), want: true, }, { - name: "node message", - message: beehivemodel.NewMessage("").SetResourceOperation("node/edge-node/default/node/edge-node", "response").SetRoute("edgecontroller", "resource"), + name: "response error message", + message: beehivemodel.NewMessage("").SetResourceOperation("node/edge-node/default/node/edge-node", "response").FillBody(fmt.Errorf("error")), want: true, }, { diff --git a/cloud/pkg/edgecontroller/controller/upstream.go b/cloud/pkg/edgecontroller/controller/upstream.go index a07880d36..07a133741 100644 --- a/cloud/pkg/edgecontroller/controller/upstream.go +++ b/cloud/pkg/edgecontroller/controller/upstream.go @@ -651,20 +651,44 @@ func kubeClientGet(uc *UpstreamController, namespace string, name string, queryT func queryInner(uc *UpstreamController, msg model.Message, queryType string) { klog.V(4).Infof("message: %s, operation is: %s, and resource is: %s", msg.GetID(), msg.GetOperation(), msg.GetResource()) - namespace, err := messagelayer.GetNamespace(msg) + var err error + var namespace, name, nodeID, resource string + namespace, err = messagelayer.GetNamespace(msg) if err != nil { klog.Warningf("message: %s process failure, get namespace failed with error: %s", msg.GetID(), err) return } - name, err := messagelayer.GetResourceName(msg) + name, err = messagelayer.GetResourceName(msg) if err != nil { klog.Warningf("message: %s process failure, get resource name failed with error: %s", msg.GetID(), err) return } - + nodeID, err = messagelayer.GetNodeID(msg) + if err != nil { + klog.Warningf("message: %s process failure, get node id failed with error: %s", msg.GetID(), err) + return + } + resource, err = messagelayer.BuildResource(nodeID, namespace, queryType, name) + if err != nil { + klog.Warningf("message: %s process failure, build message resource failed with error: %s", msg.GetID(), err) + return + } + defer func() { + if err == nil { + return + } + resMsg := model.NewMessage(msg.GetID()). + FillBody(err). + BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, resource, model.ResponseOperation) + err = uc.messageLayer.Response(*resMsg) + if err != nil { + klog.Warningf("message: %s process failure, response failed with error: %s", msg.GetID(), err) + } + }() switch msg.GetOperation() { case model.QueryOperation: - object, err := kubeClientGet(uc, namespace, name, queryType, msg) + var object metaV1.Object + object, err = kubeClientGet(uc, namespace, name, queryType, msg) if errors.IsNotFound(err) { klog.Warningf("message: %s process failure, resource not found, namespace: %s, name: %s", msg.GetID(), namespace, name) return @@ -674,24 +698,13 @@ func queryInner(uc *UpstreamController, msg model.Message, queryType string) { return } - nodeID, err := messagelayer.GetNodeID(msg) - if err != nil { - klog.Warningf("message: %s process failure, get node id failed with error: %s", msg.GetID(), err) - return - } - resource, err := messagelayer.BuildResource(nodeID, namespace, queryType, name) - if err != nil { - klog.Warningf("message: %s process failure, build message resource failed with error: %s", msg.GetID(), err) - return - } - resMsg := model.NewMessage(msg.GetID()). SetResourceVersion(object.GetResourceVersion()). FillBody(object). BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, resource, model.ResponseOperation) - err = uc.messageLayer.Response(*resMsg) - if err != nil { - klog.Warningf("message: %s process failure, response failed with error: %s", msg.GetID(), err) + rspErr := uc.messageLayer.Response(*resMsg) + if rspErr != nil { + klog.Warningf("message: %s process failure, response failed with error: %s", msg.GetID(), rspErr) return } klog.V(4).Infof("message: %s process successfully", msg.GetID()) diff --git a/edge/pkg/metamanager/client/configmap.go b/edge/pkg/metamanager/client/configmap.go index ddc329f5c..118b027a3 100644 --- a/edge/pkg/metamanager/client/configmap.go +++ b/edge/pkg/metamanager/client/configmap.go @@ -56,7 +56,10 @@ func (c *configMaps) Get(name string) (*api.ConfigMap, error) { if err != nil { return nil, fmt.Errorf("get configmap from metaManager failed, err: %v", err) } - + errContent, ok := msg.GetContent().(error) + if ok { + return nil, errContent + } content, err := msg.GetContentData() if err != nil { return nil, fmt.Errorf("parse message to configmap failed, err: %v", err) diff --git a/edge/pkg/metamanager/client/secret.go b/edge/pkg/metamanager/client/secret.go index 17ed48d42..94f031bb7 100644 --- a/edge/pkg/metamanager/client/secret.go +++ b/edge/pkg/metamanager/client/secret.go @@ -55,7 +55,10 @@ func (c *secrets) Get(name string) (*api.Secret, error) { if err != nil { return nil, fmt.Errorf("get secret from metaManager failed, err: %v", err) } - + errContent, ok := msg.GetContent().(error) + if ok { + return nil, errContent + } content, err := msg.GetContentData() if err != nil { return nil, fmt.Errorf("parse message to secret failed, err: %v", err) diff --git a/edge/pkg/metamanager/process.go b/edge/pkg/metamanager/process.go index 0ca529a37..0f347cd4e 100644 --- a/edge/pkg/metamanager/process.go +++ b/edge/pkg/metamanager/process.go @@ -43,6 +43,13 @@ func feedbackError(err error, info string, request model.Message) { } } +func feedbackResponse(message *model.Message, parentID string, resp *model.Message) { + resp.BuildHeader(resp.GetID(), parentID, resp.GetTimestamp()) + sendToEdged(resp, message.IsSync()) + respToCloud := message.NewRespByMessage(resp, OK) + sendToCloud(respToCloud) +} + func sendToEdged(message *model.Message, sync bool) { if sync { beehiveContext.SendResp(*message) @@ -408,7 +415,12 @@ func (m *metaManager) processRemoteQuery(message model.Message) { feedbackError(err, "Error to query meta in DB", message) return } - + errContent, ok := resp.GetContent().(error) + if ok { + klog.V(4).Infof("process remote query err: %v", errContent) + feedbackResponse(&message, originalID, &resp) + return + } klog.V(4).Infof("process remote query: req[%s], resp[%s]", msgDebugInfo(&message), msgDebugInfo(&resp)) content, err := resp.GetContentData() if err != nil { @@ -432,12 +444,7 @@ func (m *metaManager) processRemoteQuery(message model.Message) { if err != nil { klog.Errorf("update meta failed, %s", msgDebugInfo(&resp)) } - resp.BuildHeader(resp.GetID(), originalID, resp.GetTimestamp()) - - sendToEdged(&resp, message.IsSync()) - - respToCloud := message.NewRespByMessage(&resp, OK) - sendToCloud(respToCloud) + feedbackResponse(&message, originalID, &resp) }() } diff --git a/edge/pkg/metamanager/process_test.go b/edge/pkg/metamanager/process_test.go index 8dc74c536..59e8a6d95 100644 --- a/edge/pkg/metamanager/process_test.go +++ b/edge/pkg/metamanager/process_test.go @@ -454,6 +454,28 @@ func TestProcessQuery(t *testing.T) { } }) + //process remote query response error content + querySetterMock.EXPECT().All(gomock.Any()).Return(int64(1), errFailedDBOperation).Times(1) + querySetterMock.EXPECT().Filter(gomock.Any(), gomock.Any()).Return(querySetterMock).Times(1) + ormerMock.EXPECT().QueryTable(gomock.Any()).Return(querySetterMock).Times(1) + msg = model.NewMessage("").BuildRouter(ModuleNameEdgeHub, GroupResource, "test/"+model.ResourceTypeConfigmap, model.QueryOperation) + meta.processQuery(*msg) + message, _ = beehiveContext.Receive(ModuleNameEdgeHub) + msg = model.NewMessage(message.GetID()).BuildRouter(ModuleNameEdgeHub, GroupResource, "test/"+model.ResourceTypeConfigmap, model.QueryOperation).FillBody(fmt.Errorf("test")) + beehiveContext.SendResp(*msg) + message, _ = beehiveContext.Receive(ModuleNameEdgeHub) + msgEdged, _ := beehiveContext.Receive(ModuleNameEdged) + t.Run("ProcessRemoteQueryResponseErrorContent", func(t *testing.T) { + want := OK + if message.GetContent().(string) != want { + t.Errorf("Wrong Error message received : Wanted %v and Got %v", want, message.GetContent()) + } + wantEdged := fmt.Errorf("test").Error() + if msgEdged.GetContent().(error).Error() != wantEdged { + t.Errorf("Wrong Error message received : Wanted %v and Got %v", wantEdged, msgEdged.GetContent().(error).Error()) + } + }) + //process remote query db fail rawSetterMock.EXPECT().Exec().Return(nil, errFailedDBOperation).Times(1) ormerMock.EXPECT().Raw(gomock.Any(), gomock.Any()).Return(rawSetterMock).Times(1) |
