Aller au texte
ShemolCadre de communication cloud–edge de KubeEdge
云原生 / KubeEdge

Cadre de communication cloud–edge de KubeEdge

Pas de garantie de fraîcheur — pour être à jour, lis le source.

Références

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 :

Image de l'article
Image de l'article

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 :

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
}

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à.

go
//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 :

go
	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

InterfaceRôleImplémentation
AddModuleajouter un modulecréer un channel de messages, le mettre dans Context.channels
AddModuleGroupajouter le module à son groupchercher le channel dans channels, puis stocker group + module + channel dans typeChannels
Cleanupnettoyer un modulel’enlever de channels et typeChannels

Comm de messages beehive

Les modules enregistrés peuvent se parler. Plusieurs 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

InterfaceRôleImplémentation
Sendenvoi async vers un modulechercher le channel dans channels, y mettre le message
Receiverecevoir les messages d’un modulechercher le channel, sortir un message ; bloquer s’il n’y en a pas
SendSyncenvoi sync vers un moduleprendre 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
SendResprépondre à un syncchercher parentID dans annoChannels, y mettre le message ; s’il n’existe pas, logger
SendToGroupasync vers tous les modules d’un groupprendre tous les modules du group dans typeChannels, envoyer à chacun
SendToGroupSyncsync vers tous les modules d’un groupprendre 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.

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

Au start de Beehive, tous les modules enregistrés, puis pour chacun :

  1. Init ModuleInfo selon le type
  2. beehiveContext.AddModule
  3. beehiveContext.AddModuleGroup
  4. start() de chaque module
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 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 :

Image de l'article
Image de l'article

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.

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

Cœur :

InterfaceRôle
ServeConnserver cloud : lire en continu, transformer en message, callback de dispatch
Readlire 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
ReadMessagelire et transformer en message

Interface server et websocket

L’interface server est simple

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

L’implémentation websocket aussi

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

Le cœur de WSServer est ServeHTTP, pour l’accès des nœuds edge :

Image de l'article
Image de l'article

Interface client et websocket

Le client est simple aussi. Connect dial le server cloud et renvoie une Connection ; l’edge lit et écrit dessus.

go
//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.

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 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 :

Image de l'article
Image de l'article

Modules internes importants :

Image de l'article
Image de l'article
  • 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()

plain text
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.

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

Cœur : accès des nœuds, dispatch. Archi interne :

Image de l'article
Image de l'article

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 :

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 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 :

  1. Checks avant init, ex. limite de nœuds.
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, le mettre dans la hash de MessageDispatcher, pour le descendant.
go
//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.

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 nodeSession, le mettre dans SessionManager, le démarrer
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()

Un nodeSession par nœud — abstraction de session. SessionManager tient toutes les sessions de ce CloudHub. Au start : 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()
}

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.

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})
}

Descendant, chemin ACK :

  1. 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.
  2. 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.
  3. 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.
  4. EdgeCore sauve d’abord en store local, puis ACK. Pas d’ACK dans l’intervalle → jusqu’à 5 renvoyés. 5 échecs → drop.
  5. 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.
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

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().

go
synccontroller.Register(c.Modules.syncController)

Au start, check périodique toutes les 5 s.

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 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.

Image de l'article
Image de l'article

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.

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

Code de start :

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

Démarrage :

  1. Certs : demander à CloudCore (ou local si bien configuré), rotation, puis boucle
  2. eh.initial() crée eh.chClient, eh.chClient.Init() — viaduct ouvre la connection websocket/quic
  3. eh.pubConnectInfo(true) broadcast « connecté » aux autres modules EdgeCore
  4. Trois goroutines :
    • routeToEdge
    • routeToCloud
    • keepalive

routeToEdge : messages du cloud. Si réponse sync, beehive sendResp ; sinon, vers le group du message.

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 : messages des autres modules edge, envoi au cloud via le client websocket/quic

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 vers le cloud selon l’intervalle

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. Si erreur pendant le transfert cloud–edge, l’edge ré-init le client websocket/quic et reconnecte.