Skip to the essay
ShemolKubeEdge cloud–edge communication framework
云原生 / KubeEdge

KubeEdge cloud–edge communication framework

No freshness guaranteed — read the source if you need it current.

References

Beehive

Beehive is KubeEdge’s core messaging framework: module registration and communication between modules. CloudCore and EdgeCore both depend on it. You need Beehive’s mechanics first if you want to get KubeEdge’s design.

Beehive is a messaging framework on Go channels (the original says “goland channel”). Two core abilities: module registration, and inter-module messaging — the ModuleContext and MessageContext interfaces. Architecture:

Article image
Article image

Message format

Before Beehive’s features, the message format. A message is what modules pass around. Three parts:

  • Header:
    • ID: message ID, UUID string
    • ParentID: present if this is a response to a sync message (only then)
    • TimeStamp: when the message was made (timestamp
    • Sync: whether this is a sync-type message; true means sync
  • Route:
    • Source: where it came from
    • Group: which group it belongs to
    • Operation: the operation on the resource
    • Resource: the resource being operated on
  • Content: the payload

context data structure

Both ModuleContext and MessageContext are implemented by Context. Shape:

go
//Context is object for Context channel
type Context struct{
	//ConfigFactory goarchaius.ConfigurationFactory
	channels map[string]chan model.Message
	chsLock  sync.RwMutex
	typeChannels map[string]map[string]chan model.Message
	typeChsLock  sync.RWMutex
	anonChannels map[string]chan model.Message
	anonChsLock  sync.RWMutex
}

Full code: https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/channel/context_channel.go

  • channels — module name → message channel, for sending to that module
  • chsLock — lock for the channels map
  • typeChannels — two-level map: first key group name, second key module name, value that module’s channel
  • typeChsLock — lock for typeChannels
  • annoChannels — parentID → channel, used to send responses to sync messages
  • annoChsLock — lock for annoChannels

beehive module management

In Beehive a module is an interface. Implement it and you’re a module. Common KubeEdge modules — cloudhub, edgehub, edgeController, etc. — already do.

go
//Module Interface
type Moudule interface{
	Name() string
	Group() string
	Start()
	Enable() bool
}

Full code: https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/module.go

Module operations Beehive supports:

go
	type ModuleContext interface {
	AddModule(info *common.ModuleInfo)
	AddModuleGroup(module.group string)
	Cleanup(module string)
}

Full code: https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/context/context.go

InterfaceWhat it doesHow
AddModuleadd a modulecreate a message channel, store it in Context.channels
AddModuleGroupadd a module to its grouplook up the channel in channels, then store group + module + channel in typeChannels
Cleanupclean a moduledrop it from channels and typeChannels

beehive message communication

Registered modules can talk to each other. Several modes:

go
//MessageContext is interface for messaging syncing
type MessageContext interface {
	//async mode
	Send(module string, message model.Message)
	Receive(module string)(model.Message error)
	//sync mode
	SendSync(module string,message model.Message,timeout time.Duration)(model.Message error)
	SendResp(message model.Message)
	//group broadcast
	SendToGroup(group string,message model.Message)
	SendToGroupSync(group string,message model.Message,timeout time.Duration) error
}

https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/context/context.go

InterfaceWhat it doesHow
Sendasync send to a modulelook up the module’s channel in channels, put the message in
Receivereceive messages for a modulelook up the channel, take a message; block if none until one arrives
SendSyncsync send to a moduleget the module channel from channels, put the message in, create a new channel, add it to annoChannels keyed by messageID, wait on that channel until timeout; if a response arrives in time return it, else empty message + timeout error
SendRespsend a response to a sync messagelook up parentID in annoChannels, put the message in; if missing, log an error
SendToGroupasync send to every module in a groupget all modules in that group from typeChannels, send to each
SendToGroupSyncsync send to every module in a groupget all modules in the group, make an anonymous channel sized to the module count, send to each, wait until that many responses arrive

Registering and starting modules

When cloudcore or edgecore starts, every module is registered into the Beehive kernel. Beehive keeps a name → module map.

go
//registerModules register all the modules started in cloudcore
func registerModules(c *v1alpha1.CloudCoreConfig){
	cloudhub.Register(c.Modules.Cloudhub)
	edgecontroller.Register(c.Modules.EdgeController)
	devicecontroller.Register(c.Modules.DeviceController)
	nodeupgradejobcontroller.Register(c.Modules.NodeUpgradeJobController)
	synccontroller.Register(c.Modules.SyncController)
	cloudstream.Register(c.Modules.CloudStream,c.CommonConfig)
	router.Register(c.Modules.Router)
	dynamiccontroller.Register(c.Modules.DynamicController)
	policycontroller.Register(client.CrdConfig)
}

https://github.com/kubeedge/kubeedge/blob/master/cloud/cmd/cloudcore/app/server.go#L155-L166

When Beehive starts it gets all registered modules and, for each:

  1. Init ModuleInfo from the module type
  2. beehiveContext.AddModule
  3. beehiveContext.AddModuleGroup
  4. Call each module’s start method
go
//StartModules starts modules that are registered
func StartModules(){
	//only register channel mode,if want to use socket mode,we should also pass in common.MsgCtxTypeUSparameter
	beehiveContext.InitContext([]string{common.MsgCtxTypeChannel})

	modules := GetModules()

	for name, module := range modules{
		var m common.ModuleInfo
		switch module.contextType{
		case common.MsgCtxTypeChannel:
			m = common.ModuleInfo{
				ModuleName: name,
				ModuleType: module.contextType,
			}
		......

		default:
			klog.Exitf("unsupported context type: %s",module.contextType)
		}

		beehiveContext.AddModule(&m)
		beehiveContext.AddModuleGroup(name,module.module.Group())

		go moduleKeeper(name, module, m)
		klog.Infof("starting module %s",name)
	}
}

https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/core.go

viaduct

viaduct is the cloud–edge comms middleware. Unified abstract interfaces, server and client for different protocols, connection and data-transfer management between edge nodes and the cloud control plane. It hides protocol differences and serves the upper layer with one API. Users can pick the cloud–edge protocol by config. Built-in: websocket and quic. Later, for other edge access and business cases, viaduct can plug in new protocols quickly.

Two parts: server/client interfaces, and per-protocol implementations. The server is used by CloudCore — start protocol servers for edge access and data. The client is used by EdgeCore — dial into CloudCore. Architecture:

Article image
Article image

WebSocket as the example for server and client.

Connection interface

Connection is the core interface. viaduct supports bidirectional protocols; cloud and edge can both send. When an edge node connects, both server and client init a Connection and do full-duplex over it:

go
//the operation set of connection
type Connection interface{
	//process message from the connection
	ServeConn() //服务端从通道中持续的读取消息

	//SetReadDeadline sets the deadline for future Read calls
	//and any currently-blocked Read call.
	//A zero value for t means Read will not time out.
	SetReadDeadline(t time.Time) error

	//SetWriteDeadline sets the deadline for future write calls
	//and any currently-blocked write call.
	//Even if write times out.It may return n > 0, indicating that
	//some of the data was successfully written.
	//A zero value for t means write will not time out.
	SetWriteDeadline(t time.Time) error

	//write write raw data to the connection
	//it will open a stream for raw data
	write(raw []byte)(int,error)

	//writeMessageAsync writes data to the connection and don't care about the response
	WriteMessageAsync(msg *model.Message) error //异步

	//writeMessageSync writes data to the connection and care about the response
	WriteMessageSync(msg *model.Message)(*model.Message,error) //同步

	//ReadMessage reads message from the connection
	//it will be blocked when no message received
	//if you want to use this api for message reading
	//make sure AutoRoute be false
	ReadMessage(msg *model.Message ) error

	//RemoteAddr returns the remote network address
	RemoteAddr() net.Addr

	//LocalAddr returns the local network address
	LocalAddr() net.Addr

	//connectState return the current connection state
	ConnectionState() connectionState

	//Close closes the connection
	//Any blocked Read or Write operations will be unblocked and return errors
	Close() error
}

https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/viaduct/pkg/conn/conn.go

Core methods:

InterfaceWhat it does
ServeConncloud server: keep reading from the connection, turn into messages, callback to dispatch
Readread raw bytes
Writewrite raw bytes
WriteMessageAsyncwrite a message, don’t wait for a response
WriteMessageSyncwrite a message and wait for the peer’s response
ReadMessageread and turn into a message

Server interface and websocket

Server interface is simple

go
//protocol server
type ProtocolServer interface {
	ListenAndServerTLS() error
	close() error
}

websocket implementation is also simple

go
func (srv *WSServer) ListenAndServeTLS() error{
		return srv.server.ListenAndServeTLS("","")
}
func (srv *WSServer)Close() error{
	if srv.server != nil{
		return srv.server.Close()
	}
	return nil
}

https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/viaduct/pkg/server/ws.go

WSServer’s core is ServeHTTP — handling edge access:

Article image
Article image

Client interface and websocket

Client is also simple. Connect dials the cloud server and returns a Connection; the edge reads and writes on it.

go
//each protocol(websocket/quic) provides Connect
type Protocolclient interface{
	Connect() (conn.Connection,error)
}

The websocket client dials the cloud. On success, if Callback is set it runs, then it inits a connection object and returns it.

go
//Connect try to connect remote server
func(c *WSClient)Connect()(conn.Connection,error){
	header := c.exOpts.Header
	header.Add("ConnectionUse",string(c.options.ConnUse))
	wsConn,resp,err := c.dialer.Dial(c.options.Addr,header)
	if err ==nil{
		klog.Infof("dial %s successfully",c.options.Addr)

		//do user's processing on connection or response
		if c.exOpts.callback != nil{
			c.exOpts.Callback(wsConn,resp)
		}
		return conn.NewConnection(&conn.ConnectionOptions{
			ConnType: api.ProtocolTypeWS,
			ConnUse:  c.options.ConnUse,
			Base:     wsConn,
			Consumer: c.options.Consumer,
			Handler:  c.options.Handler,
			CtrlLane: lane.NewLane(api.ProtocolTypeWS,wsConn),
			State:    &conn.ConnectionState{
				State:  api.StatConnected,
				Headers:c.exOpts.Header.Clone(),
			},
			AutoRoute: c.options.AutoRoute,
		}),nil
	}

	//something wrong!
	var respMsg string
	if resp != nil{
		body, errRead := io.ReadAll(io.LimitReader(resp.Body,comm.MaxReadLength))
		if errRead ==nil{
			respMsg = fmt.Sprintf("response code: %d,response body: %s",resp.StatusCode,string(body))
		}else{
			respMsg = fmt.Sprintf("response code: %d",resp.StatusCode)
		}
		resp.Body.Close()
	}
	klog.Errorf("dial websocket error(%+v), response message: %s",err, respMsg)

	return nil,err
}

https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/viaduct/pkg/client/ws.go

cloudhub

Cloudhub is a CloudCore module: edge access and cloud–edge data. Middleman between Controllers and EdgeCore. It pushes downstream messages (k8s resource events, pod update, etc.) to the edge, and takes status from the edge and forwards to the right controllers. Where it sits:

Article image
Article image

Important internal pieces:

Article image
Article image
  • HTTP server: cert entry for the edge — CA, issue, rotate
  • WebSocket server: optional, WebSocket access for the edge
  • QUIC server: optional, QUIC access for the edge
  • CSI socket server: talk to the csi driver on the cloud
  • Token manager: edge access tokens, default 12h rotation
  • Certificate manager: issue and rotate edge certs
  • message handler: access management and dispatch of edge messages
  • node session manager: session lifecycle per edge node
  • message dispatcher: up and down message dispatch

Cloudhub startup

Registered when CloudCore starts; Beehive calls Start() to run the module

plain text
cloudhub.Register(c.modules.Cloudhub)

On start it first launches dispatcher.DispatchDownstream for async downstream, then inits certs — if none are configured it generates CA and server certs for later WebSocket / Quic / HTTP. Then token manager: make the edge access token and start auto-rotation. StartHTTPServer() listens mainly so EdgeCore can request certs. Then the cloudhub service itself: viaduct starts a server waiting for EdgeCore, WebSocket on TCP or QUIC on UDP. If the user needs CSI, start the CSI socket server.

go
func (ch *cloudHub) Start(){
	if !cache.WaitForCacheSync(beehiveContext.Done(),ch.informersSyncedFuncs...
	{
		klog.Errorf("unable to sync caches for objectSyncController")
		os.Exit(1)
	}

	//start dispatch message from the cloud to edge node
	go ch.dispatcher.DispatchDownstream()

	//check whether the certificates exists in the local directory.
	//and then check whether certificates exists in the secret.
	//generate if they don't exist
	if err := httpserver.PrepareAllCerts(): err!=nil{
		klog.Exit(err)
	}

	DoneTLSTunnelCerts <- true
	close(DoneTLSTunnelCerts)

	//generate Token
	if err:=httpserver.GenerateToken():err!=nil{
		klog.Exit(err)
	}

	//HttpServer mainly used to issue certificates for the edge
	go httpserver.StartHTTPServer()

	servers.StartCloudHub(ch.messageHandler)

	if hubconfig.Config.UnixSocket.Enable{
	//The uds server is only used to communicate with csi driver from kubeedge on cloud
	//It is not used to communicate  between cloud and edge
	go udsserver.StartServer(hubconfig.Config.unixSocket.Address)
	}
}

https://github.com/kubeedge/kubeedge/blob/master/cloud/pkg/cloudhub/cloudhub.go

Core jobs: edge access, and message dispatch. Internals:

Article image
Article image

Downstream send modes

Two modes for messages going to the edge. They decide how dispatch and the node session work:

ACK: after the edge gets the downstream message and saves it correctly to local store, it must ACK the cloud. If the cloud gets no ACK, it treats the message as not handled and retries until ACK.
NO-ACK: the edge does not ACK. The cloud assumes the edge got it and handled it. Messages can be lost. Usually used as the response to a sync message from the edge; if the edge doesn’t get the response, it retries.

Edge node access

Main logic is in messageHandler:

go
type Handler interface{
	//HandleConnection is invoked when a new connection arrives
	HandleConnection(connection conn.Connection)

	//HandleMessage is invoked when a new message arrives
	HandleMessage(container *mux.MessageContainer,writer mux.ResponseWriter)

	//OnEdgeNodeConnect is invoked when a new connection is established
	OnEdgeNodeConnect(info *model.HubInfo,connection conn.Connection) error
	//OnEdgeNodeDisconnect is invoked when a connection is lost
	OnEdgeNodeDisconnect(info *model.HubInfo,connection conn.Connection)

	//OnReadTransportErr is invoked when the connection read message err

}

https://github.com/kubeedge/kubeedge/blob/master/cloud/pkg/cloudhub/handler/message_handler.go

HandleConnection handles access. WebSocket example: after the WS server is up via viaduct, when an edge connects, ServeHTTP upgrades HTTP to websocket, inits a Connection, HandleConnection does:

  1. Checks before init, e.g. node count limit.
go
nodeID := connection.ConnectionState().Headers.Get("node_id")
projectID := connection.ConnectionState().Headers.Get("project_id")

if mh.SessionManager.ReachLimit(){
	klog.Errorf("Fail to serve node %s,reach node limit",nodeID)
	return
}
  1. Init nodeMessagePool, put it in MessageDispatcher’s hash, for queued downstream.
go
//init node message pool and add to the dispatcher
nodeMessagePool := common.InitNodeMessagePool(nodeID)
mh.MessageDispatcher.AddNodeMessagePool(nodeID,nodeMessagePool)

nodeMessagePool is the downstream queue. Each connecting edge gets one. Matching the two send modes, it has two queues: ACK and NO-ACK.

go
//NodeMessagePool is a collection of all downstream message sent to an
//edge node.There are two types of messages,one that requires an ack
//and another that does not.For each type of message.we use the 'queue'
//to mark the order of sending, and use the 'store' to store specific messages
type NodeMessagePool struct{
	//AckMessageStore store message that will send to edge node
	//and require acknowledgement from edge node
	AckMessageStore cache.Store
	//AckMessageQueue store message key that will send to edge node
	//and require acknowledgement from edge node
	AckMessageQueue workqueue.RateLimitingInterface
	//NoAckMessageStore store message that will send to edge node
	//and do not require acknowledgement from edge node
	NoAckMessageStore cache.Store
	//NoAckMessageQueue store message key that will send to edge node
	//and do not require acknowledgement from edge node
	NoAckMessageQueue workqueue.RateLimitingInterface
}

https://github.com/kubeedge/kubeedge/blob/master/cloud/pkg/cloudhub/common/message_pool.go

  1. Init a nodeSession, put it in SessionManager, start it
go
//create a node session for each edge node
nodeSession := session.NewNodeSession(nodeID,projectID,connection,
	keepaliveInterval,nodeMessagePool,mh.reliableClient)
//add node session  to the session manager
mh.SessionManager.AddSession(nodeSession)

//Start session for each edge node and it will keep running until
//it encounters some Transport Error from underlying connection.
nodeSession.Start()

One nodeSession per edge node — the session abstraction. SessionManager holds all sessions on this CloudHub. On start it launches the goroutines that node needs: KeepAliveCheck, SendAckMessage, SendNoAckMessage.

go
//Start the main goroutine responsible for serving node session
func (ns *NodeSession)Start(){
	klog.Infof("Start session for edge node %s",ns.nodeID)

	go ns.KeepAliveCheck()
	go ns.SendAckMessage()
	go ns.SendNoAckMessage()

	<-ns.ctx.Done()
}

Upstream / downstream dispatch

In CloudHub this is fairly simple. HandleMessage on messageHandler. viaduct parses into a MessageContainer with the message. HandleMessage does a light check, then MessageDispatcher.DispatchUpstream, off to edgeController, deviceController, etc.

go
//HandleMessage handle all the request from node
func (mh *messageHandler)HandleMessage(coantainer *mux.MessageContainer,writer mux.ResponseWriter){
	nodeID := container.Header.Get("node_id")
	projectID := container.Header.Get("project_id")

	//validate message
	if container.Message == nil{
		klog.Errorf("The message is nil for node: %s",nodeID)
		return
	}

	klog.v(4).Infof("[messageHandler]get msg from node(%s): %+v",nodeID,container.Message)

	//dispatch upstream message
	mh.MessageDispatcher.DispatchUpstream(container.Message,&model.HubInfo{ProjectID: projectID,NodeID:nodeID})
}

Downstream, ACK path:

  1. KubeEdge uses a K8s objectSync CRD to store the latest resourceVersion successfully sent to the edge. When Cloudhub starts (or restarts) it compares pending resourceVersion vs successfully sent, so it doesn’t send old messages.
  2. EdgeController, devicecontroller, etc. send to Cloudhub. MessageDispatcher routes by node name into that NodeMessagePool, and picks ACK vs NO-ACK from resource etc. While enqueueing it looks up the objectSync CR, compares versions, skips dupes.
  3. That node’s SendAckMessage goroutine pulls from NodeMessagePool in order, sends to the edge, stores the message ID on the ACK channel. When the edge ACKs, the channel fires, resourceVersion is saved on the objectSync CR, next message goes out.
  4. EdgeCore first saves to local store, then ACKs. If Cloudhub gets no ACK in the interval, it resends up to 5 times. All 5 fail → drop the event.
  5. SyncController in CloudCore handles those failures. Even if the edge got the message, the ACK can be lost on the wire. SyncController sends again to Cloudhub, downstream dispatch again, until it works.
go
func (ns *NodeSession)SendMessageWithRetry(copyMsg, msg *beehivemodel.Message)error{
	ackChan := make(chan struct{})
	ns.ackMessageCache.Store(copyMsg.GetID(),ackChan)

	//initialize retry count and timer for sending message
	retryCount := 0
	ticker := time.NewTimer(sendRetryInterval)

	err := ns.connection.WriteMessageAsync(copyMsg)
	if err !=nil{
		return err
	}

	for{
		select{
		case <- ackChan:
			ns.saveSuccessPoint(msg)
			return nil

		case <- ticker.C:
			if retryCount == 4{
				return ErrwaitTimeout
			}
			err := ns.connection.WriteMessageAsync(copyMsg)
			if err !=nil{
				return err
			}
			retryCount++
			ticker.Reset(sendRetryInterval)
		}
	}
}

https://github.com/kubeedge/kubeedge/tree/master/cloud/pkg/cloudhub/session

synccontroller

Edge networks are usually flaky, so cloud–edge drops a lot, and you can lose data. synccontroller is a CloudCore module for reliable send. KubeEdge uses objectSync to persist cloud–edge message state. While syncing, the cloud records the latest ResourceVersion successfully synced per edge node, as a CR in K8s. That keeps send order and continuity after cloud failure or edge offline restart, so you don’t resend old messages and desync. It also periodically checks and resyncs, for eventual consistency: compare each edge node’s sync state with K8s resources, push mismatches to the edge.

Registered at CloudCore start; Beehive Start() runs it.

go
synccontroller.Register(c.Modules.syncController)

On start it runs a periodic check every 5s.

go
func (sctl *SyncController) Start(){
	if !cache.WaitForCacheSync(beehiveContext.Done(),sctl.informersSyncedFuncs...){
		klog.Errorf("unable to sync caches for sync controller")
		return
	}

	sctl.deleteObjectSyncs() //check outdate sync before start to reconcile
	go wait.Until(sctl.reconcile, 5*time.Second,beehiveContext.Done())
}

https://github.com/kubeedge/kubeedge/tree/master/cloud/pkg/synccontroller

ObjectSync holds namespaced objects. Names are node name + object UUID. SyncController periodically compares the sent resourceVersion in ObjectSync with the object in K8s, then retries or deletes. When Cloudhub adds an event to NodeMessagePool it compares with what’s already there. If the pooled object is newer, drop the event; else send to the edge.

Article image
Article image

edgehub

EdgeHub is a WebSocket or QUIC client talking to CloudCore: sync cloud resource updates, report edge host and device status, etc.

Registered through Beehive when EdgeCore starts, and initialized.

go
func Register(eh *v1alpha2.EdgeHub,nodeName string){
	config.InitConfigure(eh,nodeName)
	core.Register(newEdgeHub(eh.Enable))
}

Start code:

go
func (eh *EdgeHub)Start(){
	eh.certManager = certificate.NewCertManager(config.config.EdgeHub,config.config.NodeName)
	eh.certManager.Start()
	for _, v := range GetCertSyncChannel(){
		v <- true
		close(v)
	}

	go eh.ifRotationDone()

	for{
		select{
			case <- beehiveContext.Done():
			klog.warning("EdgeHub stop")
			return
		default:
		}
		err := eh.initial()
		if err != nil{
			klog.Exitf("failed to init controller:%v",err)
			return
		}
		waitTime := time.Duration(config.Config.Heartbeat)*time.Second*2

		err = eh.chClient.Init()
		if err!= nil{
			klog.Errorf("connection failed: %v,will reconnect after %s",err,waitTime.String())
			time.Sleep(waitTime)
			continue
		}
		//execute hook func after connect
		eh.pubConnectInfo(true)
		go eh.routeToEdge()
		go eh.routeToCloud()
		go eh.keepalive()

		//wait the stop signal
		//stop authinfo manager/websocket connection
		<-eh.reconnectChan
		eh.chClient.UnInit()

		//execute hook function after disconnect
		eh.pubConnectInfo(false)

		//sleep one period of heartbeat, then try to connect cloud hub again
		klog.warningf("connection is broken, will reconnect after %s",waitTime.String())
		time.Sleep(waitTime)

	//clean channel
	clean:
		for{
			select{
			case <- eh.reconnectChan:
			default:
				break clean
			}
		}
	}
}

https://github.com/kubeedge/kubeedge/blob/master/edge/pkg/edgehub/edgehub.go

Startup:

  1. Certs: request from CloudCore (or use local if configured), start rotation, then loop
  2. eh.initial() creates eh.chClient, eh.chClient.Init() — viaduct builds the websocket/quic connection
  3. eh.pubConnectInfo(true) broadcasts “connected” to other EdgeCore modules
  4. Three goroutines:
    • routeToEdge
    • routeToCloud
    • keepalive

routeToEdge: receive messages from the cloud. If it’s a sync response, beehive sendResp; else send to the message’s group.

go
func (eh *EdgeHub)routeToEdge(){
	for{
		select{
		case<-beehiveContext.Done():
			klog.Warning("EdgeHub RouteToEdge stop")
			return
		default:
		}
		message, err := eh.chClient.Receive()
		if err!=nil{
			klog.Errorf("failed to dispatch message,discard: %v",err)
		}
	}
}

routeToCloud: receive from other edge modules, send to the cloud over the websocket/quic client

go
func (eh *EdgeHub)routeToCloud(){
	for{
		select{
		case<-beehiveContext.Done():
			klog.warning("EdgeHub RouteToCloud stop")
			return
		default:
		}
		message,err:= beehiveContext.Receive(modules.EdgeHubModuleName)
		if err !=nil{
			klog.Errorf("failed to receive message from edge: %v",err)
			time.Sleep(time.Second)
			continue
		}

		err = eh.tryThrottle(message.GetID())
		if err !=nil{
			klog.Errorf("msgID: %s,client rate limiter returned an error: %v",message.GetID(),err)
			continue
		}

		//post message to cloud hub
		err = eh.sendToCloud(message)
		if err !=nil{
			klog.Errorf("failed to send message to cloud: %v",err)
			eh.reconnectChan <- struct{}{}
			return
		}
	}
}

keepalive: heartbeat to the cloud on the heartbeat interval

go
func (eh *EdgeHub)keepalive(){
	for{
		select{
		case <-beehiveContext.Done():
			klog.warning("EdgeHub KeepAlive stop")
			return
		default:
		}
		msg := model.NewMessage("").
			BuildRouter(modules.EdgeHubModuleName,"resource","node",messagepkg.OperationKeepalive)
			FillBody("ping")

		//post message to cloud hub
		err := eh.sendToCloud(*msg)
		if err != nil{
			klog.Errorf("websocket write error: %v",err)
			eh.reconnectChan <- struct{}{}
			return
		}
		time.Sleep(time.Duration(config.Config.Heartbeat)*time.Second)
	}
}
  1. If cloud–edge transfer errors, the edge re-inits the websocket/quic client and reconnects.