diff options
| author | zhangjie <iamkadisi@163.com> | 2019-11-12 12:43:15 +0800 |
|---|---|---|
| committer | zhangjie <iamkadisi@163.com> | 2019-11-13 17:24:19 +0800 |
| commit | ef3067a429b9839c34dc2cd3421bf5ddae78180a (patch) | |
| tree | a7670b3c6601a57d670216cc03e3a7d4a53a63dc /edge | |
| parent | metamanager:uses the golang context instead of the stop channel (diff) | |
| download | kubeedge-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.go | 2 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtcontroller.go | 2 | ||||
| -rw-r--r-- | edge/pkg/edgehub/controller.go | 57 | ||||
| -rw-r--r-- | edge/pkg/edgehub/controller_test.go | 32 | ||||
| -rw-r--r-- | edge/pkg/edgehub/module.go | 73 | ||||
| -rw-r--r-- | edge/pkg/metamanager/module.go | 2 | ||||
| -rw-r--r-- | edge/pkg/metamanager/msg_processor.go | 2 | ||||
| -rw-r--r-- | edge/pkg/servicebus/servicebus.go | 151 | ||||
| -rw-r--r-- | edge/test/test.go | 6 |
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() |
