summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authorzhangjie <iamkadisi@163.com>2019-11-12 12:43:15 +0800
committerzhangjie <iamkadisi@163.com>2019-11-13 17:24:19 +0800
commitef3067a429b9839c34dc2cd3421bf5ddae78180a (patch)
treea7670b3c6601a57d670216cc03e3a7d4a53a63dc /edge
parentmetamanager:uses the golang context instead of the stop channel (diff)
downloadkubeedge-ef3067a429b9839c34dc2cd3421bf5ddae78180a.tar.gz
servivebus:uses the golang context instead of the stop channel
Signed-off-by: zhangjie <iamkadisi@163.com>
Diffstat (limited to 'edge')
-rw-r--r--edge/pkg/devicetwin/devicetwin.go2
-rw-r--r--edge/pkg/devicetwin/dtcontroller.go2
-rw-r--r--edge/pkg/edgehub/controller.go57
-rw-r--r--edge/pkg/edgehub/controller_test.go32
-rw-r--r--edge/pkg/edgehub/module.go73
-rw-r--r--edge/pkg/metamanager/module.go2
-rw-r--r--edge/pkg/metamanager/msg_processor.go2
-rw-r--r--edge/pkg/servicebus/servicebus.go151
-rw-r--r--edge/test/test.go6
9 files changed, 171 insertions, 156 deletions
diff --git a/edge/pkg/devicetwin/devicetwin.go b/edge/pkg/devicetwin/devicetwin.go
index 938fcfcd8..af057569c 100644
--- a/edge/pkg/devicetwin/devicetwin.go
+++ b/edge/pkg/devicetwin/devicetwin.go
@@ -53,7 +53,7 @@ func (dt *DeviceTwin) Start(c *beehiveContext.Context) {
klog.Errorf("Start DeviceTwin Failed, Sync Sqlite error:%v", err)
return
}
- dt.start(ctx)
+ dt.runDeviceTwin(ctx)
}
//Cleanup clean resource after quit
diff --git a/edge/pkg/devicetwin/dtcontroller.go b/edge/pkg/devicetwin/dtcontroller.go
index a3caedb6c..efc815dde 100644
--- a/edge/pkg/devicetwin/dtcontroller.go
+++ b/edge/pkg/devicetwin/dtcontroller.go
@@ -271,7 +271,7 @@ func classifyMsg(message *dttype.DTMessage) bool {
return false
}
-func (dt *DeviceTwin) start(ctx context.Context) {
+func (dt *DeviceTwin) runDeviceTwin(ctx context.Context) {
moduleNames := []string{dtcommon.MemModule, dtcommon.TwinModule, dtcommon.DeviceModule, dtcommon.CommModule}
for _, v := range moduleNames {
diff --git a/edge/pkg/edgehub/controller.go b/edge/pkg/edgehub/controller.go
index 951c53858..2e869c6e5 100644
--- a/edge/pkg/edgehub/controller.go
+++ b/edge/pkg/edgehub/controller.go
@@ -45,57 +45,6 @@ func (eh *EdgeHub) initial() (err error) {
return nil
}
-//Start will start EdgeHub
-func (eh *EdgeHub) start(ctx context.Context) {
- config.InitEdgehubConfig()
- for {
- select {
- case <-ctx.Done():
- klog.Warning("EdgeHub stop")
- return
- default:
-
- }
- err := eh.initial()
- if err != nil {
- klog.Fatalf("failed to init controller: %v", err)
- return
- }
- err = eh.chClient.Init()
- if err != nil {
- klog.Errorf("connection error, try again after 60s: %v", err)
- time.Sleep(waitConnectionPeriod)
- continue
- }
- // execute hook func after connect
- eh.pubConnectInfo(true)
- go eh.routeToEdge(ctx)
- go eh.routeToCloud(ctx)
- go eh.keepalive(ctx)
-
- // wait the stop singal
- // stop authinfo manager/websocket connection
- <-eh.retryChan
- eh.chClient.Uninit()
-
- // execute hook fun after disconnect
- eh.pubConnectInfo(false)
-
- // sleep one period of heartbeat, then try to connect cloud hub again
- time.Sleep(eh.config.HeartbeatPeriod * 2)
-
- // clean channel
- clean:
- for {
- select {
- case <-eh.retryChan:
- default:
- break clean
- }
- }
- }
-}
-
func (eh *EdgeHub) addKeepChannel(msgID string) chan model.Message {
eh.keeperLock.Lock()
defer eh.keeperLock.Unlock()
@@ -167,7 +116,7 @@ func (eh *EdgeHub) routeToEdge(ctx context.Context) {
message, err := eh.chClient.Receive()
if err != nil {
klog.Errorf("websocket read error: %v", err)
- eh.retryChan <- struct{}{}
+ eh.reconnectChan <- struct{}{}
return
}
@@ -228,7 +177,7 @@ func (eh *EdgeHub) routeToCloud(ctx context.Context) {
err = eh.sendToCloud(message)
if err != nil {
klog.Errorf("failed to send message to cloud: %v", err)
- eh.retryChan <- struct{}{}
+ eh.reconnectChan <- struct{}{}
return
}
}
@@ -251,7 +200,7 @@ func (eh *EdgeHub) keepalive(ctx context.Context) {
err := eh.sendToCloud(*msg)
if err != nil {
klog.Errorf("websocket write error: %v", err)
- eh.retryChan <- struct{}{}
+ eh.reconnectChan <- struct{}{}
return
}
diff --git a/edge/pkg/edgehub/controller_test.go b/edge/pkg/edgehub/controller_test.go
index 49d243224..928a9f04c 100644
--- a/edge/pkg/edgehub/controller_test.go
+++ b/edge/pkg/edgehub/controller_test.go
@@ -254,20 +254,20 @@ func TestRouteToEdge(t *testing.T) {
{
name: "Route to edge with proper input",
hub: &EdgeHub{
- context: beehiveContext.GetContext(beehiveContext.MsgCtxTypeChannel),
- chClient: mockAdapter,
- syncKeeper: make(map[string]chan model.Message),
- retryChan: make(chan struct{}),
+ context: beehiveContext.GetContext(beehiveContext.MsgCtxTypeChannel),
+ chClient: mockAdapter,
+ syncKeeper: make(map[string]chan model.Message),
+ reconnectChan: make(chan struct{}),
},
receiveTimes: 0,
},
{
name: "Receive Error in route to edge",
hub: &EdgeHub{
- context: beehiveContext.GetContext(beehiveContext.MsgCtxTypeChannel),
- chClient: mockAdapter,
- syncKeeper: make(map[string]chan model.Message),
- retryChan: make(chan struct{}),
+ context: beehiveContext.GetContext(beehiveContext.MsgCtxTypeChannel),
+ chClient: mockAdapter,
+ syncKeeper: make(map[string]chan model.Message),
+ reconnectChan: make(chan struct{}),
},
receiveTimes: 1,
},
@@ -278,7 +278,7 @@ func TestRouteToEdge(t *testing.T) {
mockAdapter.EXPECT().Receive().Return(*model.NewMessage("test").BuildRouter(ModuleNameEdgeHub, module.TwinGroup, "", ""), nil).Times(tt.receiveTimes)
mockAdapter.EXPECT().Receive().Return(*model.NewMessage(""), errors.New("Connection Refused")).Times(1)
go tt.hub.routeToEdge(ctx)
- stop := <-tt.hub.retryChan
+ stop := <-tt.hub.reconnectChan
if stop != struct{}{} {
t.Errorf("TestRouteToEdge error got: %v want: %v", stop, struct{}{})
}
@@ -389,9 +389,9 @@ func TestRouteToCloud(t *testing.T) {
{
name: "Route to cloud with valid input",
hub: &EdgeHub{
- context: testContext,
- chClient: mockAdapter,
- retryChan: make(chan struct{}),
+ context: testContext,
+ chClient: mockAdapter,
+ reconnectChan: make(chan struct{}),
},
},
}
@@ -404,7 +404,7 @@ func TestRouteToCloud(t *testing.T) {
testContext.AddModule(ModuleNameEdgeHub)
msg := model.NewMessage("").BuildHeader("test_id", "", 1)
testContext.Send(ModuleNameEdgeHub, *msg)
- stopChan := <-tt.hub.retryChan
+ stopChan := <-tt.hub.reconnectChan
if stopChan != struct{}{} {
t.Errorf("Error in route to cloud")
}
@@ -433,8 +433,8 @@ func TestKeepalive(t *testing.T) {
ProjectID: "foo",
NodeID: "bar",
},
- chClient: mockAdapter,
- retryChan: make(chan struct{}),
+ chClient: mockAdapter,
+ reconnectChan: make(chan struct{}),
},
},
}
@@ -448,7 +448,7 @@ func TestKeepalive(t *testing.T) {
mockAdapter.EXPECT().Send(gomock.Any()).Return(nil).Times(1)
mockAdapter.EXPECT().Send(gomock.Any()).Return(errors.New("Connection Refused")).Times(1)
go tt.hub.keepalive(ctx)
- got := <-tt.hub.retryChan
+ got := <-tt.hub.reconnectChan
if got != struct{}{} {
t.Errorf("TestKeepalive() StopChan = %v, want %v", got, struct{}{})
}
diff --git a/edge/pkg/edgehub/module.go b/edge/pkg/edgehub/module.go
index 92d888573..e8981fdcc 100644
--- a/edge/pkg/edgehub/module.go
+++ b/edge/pkg/edgehub/module.go
@@ -3,6 +3,9 @@ package edgehub
import (
"context"
"sync"
+ "time"
+
+ "k8s.io/klog"
"github.com/kubeedge/beehive/pkg/core"
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
@@ -19,21 +22,21 @@ const (
//EdgeHub defines edgehub object structure
type EdgeHub struct {
- context *beehiveContext.Context
- chClient clients.Adapter
- config *config.ControllerConfig
- retryChan chan struct{}
- cancel context.CancelFunc
- syncKeeper map[string]chan model.Message
- keeperLock sync.RWMutex
+ context *beehiveContext.Context
+ chClient clients.Adapter
+ config *config.ControllerConfig
+ reconnectChan chan struct{}
+ cancel context.CancelFunc
+ syncKeeper map[string]chan model.Message
+ keeperLock sync.RWMutex
}
// Register register edgehub
func Register() {
core.Register(&EdgeHub{
- config: &config.GetConfig().CtrConfig,
- retryChan: make(chan struct{}),
- syncKeeper: make(map[string]chan model.Message),
+ config: &config.GetConfig().CtrConfig,
+ reconnectChan: make(chan struct{}),
+ syncKeeper: make(map[string]chan model.Message),
})
}
@@ -52,7 +55,55 @@ func (eh *EdgeHub) Start(c *beehiveContext.Context) {
var ctx context.Context
eh.context = c
ctx, eh.cancel = context.WithCancel(context.Background())
- eh.start(ctx)
+
+ config.InitEdgehubConfig()
+
+ for {
+ select {
+ case <-ctx.Done():
+ klog.Warning("EdgeHub stop")
+ return
+ default:
+
+ }
+ err := eh.initial()
+ if err != nil {
+ klog.Fatalf("failed to init controller: %v", err)
+ return
+ }
+ err = eh.chClient.Init()
+ if err != nil {
+ klog.Errorf("connection error, try again after 60s: %v", err)
+ time.Sleep(waitConnectionPeriod)
+ continue
+ }
+ // execute hook func after connect
+ eh.pubConnectInfo(true)
+ go eh.routeToEdge(ctx)
+ go eh.routeToCloud(ctx)
+ go eh.keepalive(ctx)
+
+ // wait the stop singal
+ // stop authinfo manager/websocket connection
+ <-eh.reconnectChan
+ eh.chClient.Uninit()
+
+ // execute hook fun after disconnect
+ eh.pubConnectInfo(false)
+
+ // sleep one period of heartbeat, then try to connect cloud hub again
+ time.Sleep(eh.config.HeartbeatPeriod * 2)
+
+ // clean channel
+ clean:
+ for {
+ select {
+ case <-eh.reconnectChan:
+ default:
+ break clean
+ }
+ }
+ }
}
//Cleanup sets up context cleanup through Edgehub name
diff --git a/edge/pkg/metamanager/module.go b/edge/pkg/metamanager/module.go
index 19bec24e7..986ac7fec 100644
--- a/edge/pkg/metamanager/module.go
+++ b/edge/pkg/metamanager/module.go
@@ -61,7 +61,7 @@ func (m *metaManager) Start(c *beehiveContext.Context) {
}
}()
- m.mainLoop(ctx)
+ m.runMetaManager(ctx)
}
func (m *metaManager) Cleanup() {
diff --git a/edge/pkg/metamanager/msg_processor.go b/edge/pkg/metamanager/msg_processor.go
index fa4ddd106..447711bc4 100644
--- a/edge/pkg/metamanager/msg_processor.go
+++ b/edge/pkg/metamanager/msg_processor.go
@@ -628,7 +628,7 @@ func (m *metaManager) process(message model.Message) {
}
}
-func (m *metaManager) mainLoop(ctx context.Context) {
+func (m *metaManager) runMetaManager(ctx context.Context) {
go func() {
for {
select {
diff --git a/edge/pkg/servicebus/servicebus.go b/edge/pkg/servicebus/servicebus.go
index 69089c2f7..6c251ab9b 100644
--- a/edge/pkg/servicebus/servicebus.go
+++ b/edge/pkg/servicebus/servicebus.go
@@ -1,6 +1,7 @@
package servicebus
import (
+ "context"
"encoding/json"
"fmt"
"io/ioutil"
@@ -11,7 +12,7 @@ import (
"k8s.io/klog"
"github.com/kubeedge/beehive/pkg/core"
- "github.com/kubeedge/beehive/pkg/core/context"
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
"github.com/kubeedge/kubeedge/edge/pkg/servicebus/util"
@@ -24,7 +25,8 @@ const (
// servicebus struct
type servicebus struct {
- context *context.Context
+ context *beehiveContext.Context
+ cancel context.CancelFunc
}
// Register register servicebus
@@ -41,9 +43,11 @@ func (*servicebus) Group() string {
return modules.BusGroup
}
-func (sb *servicebus) Start(c *context.Context) {
+func (sb *servicebus) Start(c *beehiveContext.Context) {
// no need to call TopicInit now, we have fixed topic
+ var ctx context.Context
sb.context = c
+ ctx, sb.cancel = context.WithCancel(context.Background())
var htc = new(http.Client)
htc.Timeout = time.Second * 10
@@ -52,82 +56,93 @@ func (sb *servicebus) Start(c *context.Context) {
//Get message from channel
for {
- if msg, ok := sb.context.Receive("servicebus"); ok == nil {
- go func() {
- klog.Infof("ServiceBus receive msg")
- source := msg.GetSource()
- if source != sourceType {
- return
+ select {
+ case <-ctx.Done():
+ klog.Warning("ServiceBus stop")
+ return
+ default:
+
+ }
+ msg, err := sb.context.Receive("servicebus")
+ if err != nil {
+ klog.Warningf("servicebus receive msg error %v", err)
+ continue
+ }
+ go func() {
+ klog.Infof("ServiceBus receive msg")
+ source := msg.GetSource()
+ if source != sourceType {
+ return
+ }
+ resource := msg.GetResource()
+ r := strings.Split(resource, ":")
+ if len(r) != 2 {
+ m := "the format of resource " + resource + " is incorrect"
+ klog.Warningf(m)
+ code := http.StatusBadRequest
+ if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
+ sb.context.SendToGroup(modules.HubGroup, response)
}
- resource := msg.GetResource()
- r := strings.Split(resource, ":")
- if len(r) != 2 {
- m := "the format of resource " + resource + " is incorrect"
- klog.Warningf(m)
- code := http.StatusBadRequest
- if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
- sb.context.SendToGroup(modules.HubGroup, response)
- }
- return
+ return
+ }
+ content, err := json.Marshal(msg.GetContent())
+ if err != nil {
+ klog.Errorf("marshall message content failed %v", err)
+ m := "error to marshal request msg content"
+ code := http.StatusBadRequest
+ if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
+ sb.context.SendToGroup(modules.HubGroup, response)
}
- content, err := json.Marshal(msg.GetContent())
- if err != nil {
- klog.Errorf("marshall message content failed %v", err)
- m := "error to marshal request msg content"
- code := http.StatusBadRequest
- if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
- sb.context.SendToGroup(modules.HubGroup, response)
- }
- return
+ return
+ }
+ var httpRequest util.HTTPRequest
+ if err := json.Unmarshal(content, &httpRequest); err != nil {
+ m := "error to parse http request"
+ code := http.StatusBadRequest
+ klog.Errorf(m, err)
+ if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
+ sb.context.SendToGroup(modules.HubGroup, response)
}
- var httpRequest util.HTTPRequest
- if err := json.Unmarshal(content, &httpRequest); err != nil {
- m := "error to parse http request"
- code := http.StatusBadRequest
- klog.Errorf(m, err)
- if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
- sb.context.SendToGroup(modules.HubGroup, response)
- }
- return
+ return
+ }
+ operation := msg.GetOperation()
+ targetURL := "http://127.0.0.1:" + r[0] + "/" + r[1]
+ resp, err := uc.HTTPDo(operation, targetURL, httpRequest.Header, httpRequest.Body)
+ if err != nil {
+ m := "error to call service"
+ code := http.StatusNotFound
+ klog.Errorf(m, err)
+ if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
+ sb.context.SendToGroup(modules.HubGroup, response)
}
- operation := msg.GetOperation()
- targetURL := "http://127.0.0.1:" + r[0] + "/" + r[1]
- resp, err := uc.HTTPDo(operation, targetURL, httpRequest.Header, httpRequest.Body)
- if err != nil {
- m := "error to call service"
- code := http.StatusNotFound
- klog.Errorf(m, err)
- if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
- sb.context.SendToGroup(modules.HubGroup, response)
- }
- return
+ return
+ }
+ resp.Body = http.MaxBytesReader(nil, resp.Body, maxBodySize)
+ resBody, err := ioutil.ReadAll(resp.Body)
+ if err != nil {
+ if err.Error() == "http: request body too large" {
+ err = fmt.Errorf("response body too large")
}
- resp.Body = http.MaxBytesReader(nil, resp.Body, maxBodySize)
- resBody, err := ioutil.ReadAll(resp.Body)
- if err != nil {
- if err.Error() == "http: request body too large" {
- err = fmt.Errorf("response body too large")
- }
- m := "error to receive response, err: " + err.Error()
- code := http.StatusInternalServerError
- klog.Errorf(m, err)
- if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
- sb.context.SendToGroup(modules.HubGroup, response)
- }
- return
+ m := "error to receive response, err: " + err.Error()
+ code := http.StatusInternalServerError
+ klog.Errorf(m, err)
+ if response, err := buildErrorResponse(msg.GetID(), m, code); err == nil {
+ sb.context.SendToGroup(modules.HubGroup, response)
}
+ return
+ }
- response := util.HTTPResponse{Header: resp.Header, StatusCode: resp.StatusCode, Body: resBody}
- responseMsg := model.NewMessage(msg.GetID())
- responseMsg.Content = response
- responseMsg.SetRoute("servicebus", modules.UserGroup)
- sb.context.SendToGroup(modules.HubGroup, *responseMsg)
- }()
- }
+ response := util.HTTPResponse{Header: resp.Header, StatusCode: resp.StatusCode, Body: resBody}
+ responseMsg := model.NewMessage(msg.GetID())
+ responseMsg.Content = response
+ responseMsg.SetRoute("servicebus", modules.UserGroup)
+ sb.context.SendToGroup(modules.HubGroup, *responseMsg)
+ }()
}
}
func (sb *servicebus) Cleanup() {
+ sb.cancel()
sb.context.Cleanup(sb.Name())
}
diff --git a/edge/test/test.go b/edge/test/test.go
index e97cb662a..1cd494eb4 100644
--- a/edge/test/test.go
+++ b/edge/test/test.go
@@ -12,7 +12,7 @@ import (
"k8s.io/klog"
"github.com/kubeedge/beehive/pkg/core"
- "github.com/kubeedge/beehive/pkg/core/context"
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
@@ -29,7 +29,7 @@ func Register() {
}
type testManager struct {
- context *context.Context
+ context *beehiveContext.Context
moduleWait *sync.WaitGroup
}
@@ -217,7 +217,7 @@ func (tm *testManager) configmapHandler(w http.ResponseWriter, req *http.Request
}
}
-func (tm *testManager) Start(c *context.Context) {
+func (tm *testManager) Start(c *beehiveContext.Context) {
tm.context = c
defer tm.Cleanup()