Pas de garantie de fraîcheur — pour être à jour, lis le source.
Références
- Cours ouvert KubeEdge edge computing cloud native 15 — arborescence du projet et source du cadre de comm cloud–edge
- Série d’analyse du source kubeedge (2) : cloudhub
Beehive
Beehive est le cadre de messages au cœur de KubeEdge : enregistrement des modules et comm entre modules. CloudCore et EdgeCore en dépendent. Il faut d’abord Beehive pour comprendre le design de KubeEdge.
Beehive est un cadre de messages sur des channels Go (l’original dit « goland channel »). Deux capacités : gestion d’enregistrement des modules, gestion de la comm entre modules — interfaces ModuleContext et MessageContext. Archi :

Format de message
Avant les fonctions de Beehive, le format. Un message, c’est le porteur entre modules. Trois parties :
- Header :
- ID : id du message, UUID
- ParentID : présent si c’est une réponse à un message sync (seulement dans ce cas)
- TimeStamp : moment de création (timestamp
- Sync : est-ce un message sync ;
true= sync
- Route :
- Source : d’où ça vient
- Group : le group
- Operation : l’opération sur la ressource
- Resource : la ressource
- Content : le contenu
Structure de context
ModuleContext et MessageContext sont implémentés par Context. Forme :
//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
}Code complet : https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/channel/context_channel.go
- channels — nom de module → channel de messages, pour envoyer à ce module
- chsLock — lock de channels
- typeChannels — map à deux niveaux : 1er clé nom de group, 2e clé nom de module, value le channel du module
- typeChsLock — lock de typeChannels
- annoChannels — parentID → channel, pour envoyer la réponse d’un sync
- annoChsLock — lock de annoChannels
Gestion des modules beehive
Dans Beehive un module est une interface. Tu l’implémentes, tu es un module. Les modules KubeEdge courants — cloudhub, edgehub, edgeController, etc. — l’ont déjà.
//Module Interface
type Moudule interface{
Name() string
Group() string
Start()
Enable() bool
}Code complet : https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/module.go
Opérations module supportées :
type ModuleContext interface {
AddModule(info *common.ModuleInfo)
AddModuleGroup(module.group string)
Cleanup(module string)
}Code complet : https://github.com/kubeedge/kubeedge/blob/master/staging/src/github.com/kubeedge/beehive/pkg/core/context/context.go
| Interface | Rôle | Implémentation |
|---|---|---|
| AddModule | ajouter un module | créer un channel de messages, le mettre dans Context.channels |
| AddModuleGroup | ajouter le module à son group | chercher le channel dans channels, puis stocker group + module + channel dans typeChannels |
| Cleanup | nettoyer un module | l’enlever de channels et typeChannels |
Comm de messages beehive
Les modules enregistrés peuvent se parler. Plusieurs modes :
//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
}| Interface | Rôle | Implémentation |
|---|---|---|
| Send | envoi async vers un module | chercher le channel dans channels, y mettre le message |
| Receive | recevoir les messages d’un module | chercher le channel, sortir un message ; bloquer s’il n’y en a pas |
| SendSync | envoi sync vers un module | prendre le channel, y mettre le message, créer un nouveau channel, l’ajouter à annoChannels sous messageID, attendre dessus jusqu’au timeout ; réponse à temps → la renvoyer, sinon message vide + erreur timeout |
| SendResp | répondre à un sync | chercher parentID dans annoChannels, y mettre le message ; s’il n’existe pas, logger |
| SendToGroup | async vers tous les modules d’un group | prendre tous les modules du group dans typeChannels, envoyer à chacun |
| SendToGroupSync | sync vers tous les modules d’un group | prendre tous les modules, créer un channel anonyme de taille = nombre de modules, send, attendre que le nombre de réponses = size |
Enregistrement et démarrage
Au start de cloudcore ou edgecore, tous les modules sont enregistrés dans le noyau Beehive. Map nom → module.
//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
Au start de Beehive, tous les modules enregistrés, puis pour chacun :
- Init ModuleInfo selon le type
- beehiveContext.AddModule
- beehiveContext.AddModuleGroup
- start() de chaque module
//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)
}
}viaduct
viaduct est le middleware de comm cloud–edge. Interfaces abstraites unifiées, serveur et client par protocole, gestion des connexions et du transfert entre nœuds edge et plan de contrôle cloud. Il cache les différences de protocoles, une API vers le haut. L’utilisateur choisit le protocole par config. Built-in : websocket et quic. Ensuite, selon les scènes d’accès edge, viaduct peut brancher un nouveau protocole vite.
Deux parties : interfaces serveur/client, et implémentations par protocole. Le serveur est utilisé par CloudCore — démarrer des servers de protocoles pour l’accès edge et le transfert. Le client est utilisé par EdgeCore — se connecter à CloudCore. Archi :

Exemple WebSocket pour serveur et client.
Interface Connection
Connection est l’interface cœur. viaduct supporte le bidirectionnel ; cloud et edge envoient dans les deux sens. Quand un nœud edge se connecte, server et client initient une Connection et font du full-duplex dessus.
//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
}Cœur :
| Interface | Rôle |
|---|---|
| ServeConn | server cloud : lire en continu, transformer en message, callback de dispatch |
| Read | lire des bytes bruts |
| Write | écrire des bytes bruts |
| WriteMessageAsync | écrire un message, sans attendre de réponse |
| WriteMessageSync | écrire un message et attendre la réponse du pair |
| ReadMessage | lire et transformer en message |
Interface server et websocket
L’interface server est simple
//protocol server
type ProtocolServer interface {
ListenAndServerTLS() error
close() error
}L’implémentation websocket aussi
func (srv *WSServer) ListenAndServeTLS() error{
return srv.server.ListenAndServeTLS("","")
}
func (srv *WSServer)Close() error{
if srv.server != nil{
return srv.server.Close()
}
return nil
}Le cœur de WSServer est ServeHTTP, pour l’accès des nœuds edge :

Interface client et websocket
Le client est simple aussi. Connect dial le server cloud et renvoie une Connection ; l’edge lit et écrit dessus.
//each protocol(websocket/quic) provides Connect
type Protocolclient interface{
Connect() (conn.Connection,error)
}Le client websocket dial le cloud. Si ça marche, si Callback n’est pas vide il le lance, puis init la connection et la renvoie.
//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
}cloudhub
Cloudhub est un module de CloudCore : accès des nœuds edge et transfert cloud–edge. Intermédiaire entre Controllers et EdgeCore. Il pousse les messages descendants (événements de ressources k8s, pod update, etc.) vers l’edge, et prend les statuts de l’edge vers les bons controllers. Place dans KubeEdge :

Modules internes importants :

- HTTP server : entrée certs pour l’edge — CA, émission, rotation
- WebSocket server : activable, accès WebSocket pour l’edge
- QUIC server : activable, accès QUIC pour l’edge
- CSI socket server : parler au csi driver côté cloud
- Token manager : tokens d’accès edge, rotation 12h par défaut
- Certificate manager : émission et rotation des certs edge
- message handler : accès et dispatch des messages edge
- node session manager : cycle de vie des sessions
- message dispatcher : dispatch montée / descente
Démarrage Cloudhub
Enregistré au start de CloudCore ; Beehive appelle Start()
cloudhub.Register(c.modules.Cloudhub)Au start : d’abord dispatcher.DispatchDownstream en goroutine pour le descendant async, puis init des certs — s’il n’y en a pas, génération auto CA + certs serveur pour WebSocket / Quic / HTTP. Puis token manager : token d’accès edge et rotation auto. StartHTTPServer() écoute surtout pour qu’EdgeCore demande des certs. Puis le service cloudhub : viaduct démarre un server, attend EdgeCore, WebSocket sur TCP ou QUIC sur UDP. Si CSI, CSI socket server aussi.
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
Cœur : accès des nœuds, dispatch. Archi interne :

Modes d’envoi descendant
Deux modes, qui décident le dispatch et le traitement de session :
ACK : l’edge, après avoir reçu le descendant et l’avoir bien sauvé en store local, envoie un ACK au cloud. Pas d’ACK = pas traité, retry jusqu’à ACK.
NO-ACK : pas d’ACK. Le cloud considère que l’edge a reçu et traité. Le message peut se perdre. Souvent pour répondre à un sync de l’edge ; si l’edge n’a pas la réponse, il retry.
Accès d’un nœud edge
Logique principale dans messageHandler :
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 gère l’accès. Exemple WebSocket : après le WS server via viaduct, quand un edge arrive, ServeHTTP upgrade HTTP en websocket, init Connection, HandleConnection :
- Checks avant init, ex. limite de nœuds.
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
}- Init nodeMessagePool, le mettre dans la hash de MessageDispatcher, pour le descendant.
//init node message pool and add to the dispatcher
nodeMessagePool := common.InitNodeMessagePool(nodeID)
mh.MessageDispatcher.AddNodeMessagePool(nodeID,nodeMessagePool)nodeMessagePool = file de messages descendants. Un par nœud à la connexion. Deux files, ACK et NO-ACK.
//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
- Init nodeSession, le mettre dans SessionManager, le démarrer
//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()Un nodeSession par nœud — abstraction de session. SessionManager tient toutes les sessions de ce CloudHub. Au start : KeepAliveCheck, SendAckMessage, SendNoAckMessage.
//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()
}Dispatch montée / descente
Dans CloudHub c’est assez simple. HandleMessage. viaduct parse en MessageContainer avec le message. HandleMessage vérifie un peu, puis DispatchUpstream vers edgeController, deviceController, etc.
//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})
}Descendant, chemin ACK :
- KubeEdge utilise un CRD objectSync K8s pour le dernier resourceVersion bien envoyé à l’edge. Au start / restart de Cloudhub, il compare resourceVersion à envoyer vs déjà envoyé, pour ne pas renvoyer du vieux.
- EdgeController, devicecontroller, etc. envoient à Cloudhub. MessageDispatcher route par nom de nœud vers le NodeMessagePool, et choisit le mode selon resource etc. En enqueue, il lit l’objectSync CR, compare les versions, évite les doublons.
- SendAckMessage du nœud tire dans l’ordre, envoie à l’edge, stocke l’id sur le channel ACK. ACK de l’edge → notification, resourceVersion sauvé sur l’objectSync CR, message suivant.
- EdgeCore sauve d’abord en store local, puis ACK. Pas d’ACK dans l’intervalle → jusqu’à 5 renvoyés. 5 échecs → drop.
- SyncController dans CloudCore gère ces échecs. Même si l’edge a le message, l’ACK peut se perdre. SyncController renvoie à Cloudhub, dispatch descendant encore, jusqu’à succès.
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
En edge computing le réseau est souvent instable, donc beaucoup de coupures, risque de perdre des data. synccontroller est un module CloudCore pour l’envoi fiable. KubeEdge persist l’état des messages cloud–edge avec objectSync. Pendant la sync, le cloud enregistre le dernier ResourceVersion réussi par nœud, en CR dans K8s. Ça garde l’ordre et la continuité après panne cloud ou restart edge offline, sans renvoyer du vieux et désync. Il vérifie aussi périodiquement : état de sync de chaque nœud vs ressources K8s, pousse les écarts, cohérence à terme.
Enregistré au start de CloudCore ; Beehive Start().
synccontroller.Register(c.Modules.syncController)Au start, check périodique toutes les 5 s.
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 tient des objets namespaced. Noms = nom de nœud + UUID. SyncController compare périodiquement le resourceVersion envoyé dans ObjectSync avec l’objet K8s, puis retry ou delete. Quand Cloudhub ajoute un événement au NodeMessagePool, il compare. Si l’objet dans la pool est plus neuf, drop ; sinon envoi à l’edge.

edgehub
EdgeHub est un client WebSocket ou QUIC vers CloudCore : sync des updates de ressources cloud, report d’état hôte et devices edge, etc.
Enregistré via Beehive au start d’EdgeCore, puis init.
func Register(eh *v1alpha2.EdgeHub,nodeName string){
config.InitConfigure(eh,nodeName)
core.Register(newEdgeHub(eh.Enable))
}Code de start :
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
Démarrage :
- Certs : demander à CloudCore (ou local si bien configuré), rotation, puis boucle
- eh.initial() crée eh.chClient, eh.chClient.Init() — viaduct ouvre la connection websocket/quic
- eh.pubConnectInfo(true) broadcast « connecté » aux autres modules EdgeCore
- Trois goroutines :
- routeToEdge
- routeToCloud
- keepalive
routeToEdge : messages du cloud. Si réponse sync, beehive sendResp ; sinon, vers le group du message.
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 : messages des autres modules edge, envoi au cloud via le client websocket/quic
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 vers le cloud selon l’intervalle
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)
}
}- Si erreur pendant le transfert cloud–edge, l’edge ré-init le client websocket/quic et reconnecte.