summaryrefslogtreecommitdiff
path: root/tests
diff options
context:
space:
mode:
authorElaine4CY <chenyuncy2017@outlook.com>2021-03-02 16:19:32 +0800
committerWintonChan <cw_hao@163.com>2021-05-31 10:09:49 +0800
commitfec89d4cbaa7d4739ef1bf4bfc86adb79153cd60 (patch)
treeea3af3bef498703a2f0fbb6846f1b72375effe1e /tests
parentadd admission of servicebus in rule and rule endpoint (diff)
downloadkubeedge-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.go49
-rw-r--r--tests/e2e/utils/common.go18
-rw-r--r--tests/e2e/utils/rule.go55
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)