diff options
| author | Elaine4CY <chenyuncy2017@outlook.com> | 2021-03-02 16:19:32 +0800 |
|---|---|---|
| committer | WintonChan <cw_hao@163.com> | 2021-05-31 10:09:49 +0800 |
| commit | fec89d4cbaa7d4739ef1bf4bfc86adb79153cd60 (patch) | |
| tree | ea3af3bef498703a2f0fbb6846f1b72375effe1e /tests | |
| parent | add admission of servicebus in rule and rule endpoint (diff) | |
| download | kubeedge-fec89d4cbaa7d4739ef1bf4bfc86adb79153cd60.tar.gz | |
add e2e test: Rest to ServiceBus
Signed-off-by: Elaine4CY <chenyuncy2017@outlook.com>
Diffstat (limited to 'tests')
| -rw-r--r-- | tests/e2e/deployment/rule_crd_test.go | 49 | ||||
| -rw-r--r-- | tests/e2e/utils/common.go | 18 | ||||
| -rw-r--r-- | tests/e2e/utils/rule.go | 55 |
3 files changed, 115 insertions, 7 deletions
diff --git a/tests/e2e/deployment/rule_crd_test.go b/tests/e2e/deployment/rule_crd_test.go index 4b29036dc..dc8ba0d92 100644 --- a/tests/e2e/deployment/rule_crd_test.go +++ b/tests/e2e/deployment/rule_crd_test.go @@ -2,11 +2,10 @@ package deployment import ( "bytes" - "net/http" - "time" - . "github.com/onsi/ginkgo" . "github.com/onsi/gomega" + "net/http" + "time" v1 "github.com/kubeedge/kubeedge/cloud/pkg/apis/rules/v1" "github.com/kubeedge/kubeedge/tests/e2e/utils" @@ -99,7 +98,7 @@ var _ = Describe("Rule Management test in E2E scenario", func() { }() time.Sleep(3 * time.Second) // call rest api to send message to edge. - IsSend, statusCode := utils.SendMsg("http://127.0.0.1:9443/edge-node/default/ccc", []byte(msg)) + IsSend, statusCode := utils.SendMsg("http://127.0.0.1:9443/edge-node/default/ccc", []byte(msg), nil) Expect(IsSend).Should(BeTrue()) Expect(statusCode).Should(Equal(http.StatusOK)) Eventually(func() bool { @@ -136,5 +135,47 @@ var _ = Describe("Rule Management test in E2E scenario", func() { return b.String() == msg }, "30s", "2s").Should(Equal(true), "endpoint not listen any request.") }) + It("E2E_CREATE_RULE_3: Create rule: rest to servicebus.", func() { + var ruleList v1.RuleList + // create rest ruleendpoint + IsRestRuleEndpointCreated, status := utils.HandleRuleEndpoint(http.MethodPost, ctx.Cfg.K8SMasterForKubeEdge+RuleEndpointHandler, "", utils.RestType) + Expect(IsRestRuleEndpointCreated).Should(BeTrue()) + Expect(status).Should(Equal(http.StatusCreated)) + // create servicebus ruleendpoint + IsServicebusRuleEndpointCreated, status := utils.HandleRuleEndpoint(http.MethodPost, ctx.Cfg.K8SMasterForKubeEdge+RuleEndpointHandler, "", utils.ServicebusType) + Expect(IsServicebusRuleEndpointCreated).Should(BeTrue()) + Expect(status).Should(Equal(http.StatusCreated)) + // create rule: rest to servicebus + IsRuleCreated, statusCode := utils.HandleRule(http.MethodPost, ctx.Cfg.K8SMasterForKubeEdge+RuleHandler, "", utils.RestType, utils.ServicebusType) + Expect(IsRuleCreated).Should(BeTrue()) + Expect(statusCode).Should(Equal(http.StatusCreated)) + newRule := utils.NewRule(utils.RestType, utils.ServicebusType) + _, err := utils.GetRuleList(&ruleList, ctx.Cfg.K8SMasterForKubeEdge+RuleHandler, newRule) + Expect(err).To(BeNil()) + msg := "Hello World!" + msgHeader := map[string]string{ + "user": "I am user", + "passwd": "I am passwd", + } + b := new(bytes.Buffer) + go func() { + recieveMsg, err := utils.StartEchoServer() + if err != nil { + utils.Fatalf("fail to call edge-app's API. reason: %s. ", err.Error()) + } + b.WriteString(recieveMsg) + }() + time.Sleep(3 * time.Second) + // call rest api to send message to edge. + IsSend, statusCode := utils.SendMsg("http://127.0.0.1:9443/edge-node/default/ddd", []byte(msg), msgHeader) + Expect(IsSend).Should(BeTrue()) + Expect(statusCode).Should(Equal(http.StatusOK)) + Eventually(func() bool { + utils.Infof("receive: %s, sent msg: %s ", b.String(), msg) + newMsg := "Reply from server: " + msg + " Header of the message: [user]: " + msgHeader["user"] + + ", [passwd]: " + msgHeader["passwd"] + return b.String() == newMsg + }, "30s", "2s").Should(Equal(true), "servicebus did not return any response.") + }) }) }) diff --git a/tests/e2e/utils/common.go b/tests/e2e/utils/common.go index 7ce252ddc..77a286566 100644 --- a/tests/e2e/utils/common.go +++ b/tests/e2e/utils/common.go @@ -895,7 +895,7 @@ func CompareTwin(deviceTwin map[string]*MsgTwin, expectedDeviceTwin map[string]* return true } -func SendMsg(url string, message []byte) (bool, int) { +func SendMsg(url string, message []byte, header map[string]string) (bool, int) { var req *http.Request var err error @@ -911,6 +911,9 @@ func SendMsg(url string, message []byte) (bool, int) { Fatalf("Frame HTTP request failed, request: %s, reason: %v", req.URL.String(), err) return false, 0 } + for k, v := range header { + req.Header.Add(k, v) + } t := time.Now() resp, err := client.Do(req) if err != nil { @@ -931,8 +934,21 @@ func StartEchoServer() (string, error) { Errorf("Echo server write failed. reason: %s", err.Error()) } } + url := func(response http.ResponseWriter, request *http.Request) { + b, _ := ioutil.ReadAll(request.Body) + var buff bytes.Buffer + buff.WriteString("Reply from server: ") + buff.Write(b) + buff.WriteString(" Header of the message: [user]: " + request.Header.Get("user") + + ", [passwd]: " + request.Header.Get("passwd")) + if _, err := response.Write(buff.Bytes()); err != nil { + Errorf("Echo server write failed. reason: %s", err.Error()) + } + r <- buff.String() + } mux := http.NewServeMux() mux.HandleFunc("/echo", echo) + mux.HandleFunc("/url", url) server := &http.Server{Addr: "0.0.0.0:9000", Handler: mux} go func() { err := server.ListenAndServe() diff --git a/tests/e2e/utils/rule.go b/tests/e2e/utils/rule.go index e48c1a088..eba8418ab 100644 --- a/tests/e2e/utils/rule.go +++ b/tests/e2e/utils/rule.go @@ -17,8 +17,9 @@ import ( ) const ( - RestType = "rest" - EventbusType = "eventbus" + RestType = "rest" + EventbusType = "eventbus" + ServicebusType = "servicebus" ) func NewRule(sourceType, targetType string) *rulesv1.Rule { @@ -27,6 +28,8 @@ func NewRule(sourceType, targetType string) *rulesv1.Rule { return NewRest2EventbusRule() case sourceType == EventbusType && targetType == RestType: return NewEventbus2RestRule() + case sourceType == RestType && targetType == ServicebusType: + return NewRest2ServicebusRule() } return nil } @@ -86,12 +89,41 @@ func NewRest2EventbusRule() *rulesv1.Rule { return &rule } +func NewRest2ServicebusRule() *rulesv1.Rule { + rule := rulesv1.Rule{ + TypeMeta: v1.TypeMeta{ + Kind: "Rule", + APIVersion: "rules.kubeedge.io/v1", + }, + ObjectMeta: v1.ObjectMeta{ + Name: "rule-rest-servicebus-test", + Namespace: Namespace, + }, + Spec: rulesv1.RuleSpec{ + Source: "rest-test", + SourceResource: map[string]string{ + "path": "/ddd", + }, + Target: "servicebus-test", + TargetResource: map[string]string{ + "path": "/url", + }, + }, + Status: rulesv1.RuleStatus{ + Errors: []string{}, + }, + } + return &rule +} + func NewRuleEndpoint(endpointType string) *rulesv1.RuleEndpoint { switch endpointType { case RestType: return newRestRuleEndpoint() case EventbusType: return newEventBusRuleEndpoint() + case ServicebusType: + return newServiceBusRuleEndpoint() } return newRestRuleEndpoint() } @@ -130,6 +162,25 @@ func newEventBusRuleEndpoint() *rulesv1.RuleEndpoint { return &eventbusRuleEndpoint } +func newServiceBusRuleEndpoint() *rulesv1.RuleEndpoint { + servicebusRuleEndpoint := rulesv1.RuleEndpoint{ + TypeMeta: v1.TypeMeta{ + Kind: "RuleEndpoint", + APIVersion: "rules.kubeedge.io/v1", + }, + ObjectMeta: v1.ObjectMeta{ + Name: "servicebus-test", + Namespace: Namespace, + }, + Spec: rulesv1.RuleEndpointSpec{ + RuleEndpointType: ServicebusType, + Properties: map[string]string{ + "service_port": "9000"}, + }, + } + return &servicebusRuleEndpoint +} + // GetRuleList to get the rule list and verify whether the contents of the rule matches with what is expected func GetRuleList(list *rulesv1.RuleList, getRuleAPI string, expectedRule *rulesv1.Rule) ([]rulesv1.Rule, error) { resp, err := SendHTTPRequest(http.MethodGet, getRuleAPI) |
