鮮度は保証しない。確実にするならソースを読んでほしい。
参考資料
Beehive
Beehive は KubeEdge の中核メッセージ通信フレームワークで、モジュールの登録とモジュール間通信に使う。CloudCore と EdgeCore はどちらも beehive に依存する。だから先に Beehive の仕組みを知らないと、KubeEdge の設計は分からない。
Beehive は Go channel(原文は goland channel)ベースのメッセージ通信。核心は二つ:module 登録管理と module 間通信。beehive の ModuleContext と MessageContext。構成は下図:

message の通信フォーマット
beehive の機能の前に、メッセージの形。message は module 間の運び屋で、三つの部分がある:
- Header:
- ID:メッセージ ID、UUID 文字列
- ParentID:同期メッセージへの応答なら parentID がある(同期応答にしかない)
- TimeStamp:生成時刻(タイムスタンプ
- Sync:同期タイプかどうか。
trueなら同期メッセージ
- Route:
- Source:出所
- Group:所属 group
- Operation:リソースへの操作
- Resource:操作するリソース
- Content:中身
context のデータ構造
ModuleContext と MessageContext はどちらも Context が実装する。形はこう:
//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
}- channels — モジュール名と対応するメッセージ channel。そのモジュールへ送るため
- chsLock — channels map のロック
- typeChannels — 二段 map。一段目は group 名、二段目は module 名、value はその module の channel
- typeChsLock — typeChannels のロック
- annoChannels — parentID から channel へのマップ。同期メッセージの応答を送る
- annoChsLock — annoChannels のロック
beehive の module 管理
beehive では module はインタフェース。実装すれば module。KubeEdge の cloudhub、edgehub、edgeController などはすでに実装している。
//Module Interface
type Moudule interface{
Name() string
Group() string
Start()
Enable() bool
}Beehive が支える module 操作:
type ModuleContext interface {
AddModule(info *common.ModuleInfo)
AddModuleGroup(module.group string)
Cleanup(module string)
}| インタフェース | 機能 | 実装 |
|---|---|---|
| AddModule | module を追加 | message 型の channel を作り、Context の channels map に入れる |
| AddModuleGroup | module を所属 group へ | channels map から channel を引き、group と module と channel を typeChannels に入れる |
| Cleanup | module を掃除 | channels と typeChannels から消す |
beehive のメッセージ通信
登録されたモジュール同士は通信できる。やり方は複数:
//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
}| インタフェース | 機能 | 実装 |
|---|---|---|
| Send | 指定 module へ非同期送信 | channels map からその module の channel を引き、メッセージを入れる |
| Receive | 指定 module 宛てのメッセージを受信 | Channels map から channel を引き、取り出す。なければ届くまでブロック |
| SendSync | 指定 module へ同期送信 | channels map から channel を取りメッセージを入れ、新しい channel を作って annoChannels に messageID をキーに追加し、その channel で応答を待つ。タイムアウト前に来ればそれを返し、来なければ空メッセージとタイムアウトエラー |
| SendResp | 同期メッセージへの応答 | parentID で annoChannels を引き、メッセージを入れる。なければエラーを記録 |
| SendToGroup | 指定 group の全 module へ非同期 | typeChannels からその group の全 module を取り、順に送る |
| SendToGroupSync | 指定 group の全 module へ同期 | typeChannels から全 module を取り、module 数と同じサイズの匿名 channel を作り、send して、届いた数が size になるまで待つ |
beehive の登録と起動
cloudcore か edgecore の起動時、全 module を beehive カーネルに登録する。名前から 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
beehive 起動時、登録済み module を全部取り、順に:
- module の種類から moduleInfo を初期化
- beehiveContext.AddModule
- beehiveContext.AddModuleGroup
- 各 module の start を呼ぶ
//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 は KubeEdge 雲辺通信のミドルウェア。統一した抽象インタフェースの上に、プロトコルごとのサーバとクライアントを提供し、エッジノードとクラウド管理面の接続とデータ転送を扱う。プロトコルの差を隠し、上層には一つのインタフェース。設定で雲辺のプロトコルを選べる。今は websocket と quic が内蔵。あとのエッジ接続や業務に合わせて、viaduct で新しいプロトコルをすぐ足せる。
主に二部:サーバ/クライアントのインタフェースと、各プロトコルの実装。サーバは CloudCore が使い、プロトコルごとの server を起動してエッジの接入と転送。クライアントは EdgeCore が使い、CloudCore へつなぐ。構成:

WebSocket を例に、サーバとクライアントを見る。
Connection インタフェース
Connection が中核。viaduct は双方向プロトコルで、クラウドとエッジは双方向に送れる。エッジが接入するとき、server も client も Connection を初期化し、これで全二重する。
//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
}中核:
| インタフェース | 機能 |
|---|---|
| ServeConn | クラウド server が使う。connection から読み続けて message にし、コールバックで配送 |
| Read | 生バイトを読む |
| Write | 生バイトを書く |
| WriteMessageAsync | message を書き、応答は気にしない |
| WriteMessageSync | message を書き、対向の応答を待つ |
| ReadMessage | 読んで message にする |
server インタフェースと websocket
server のインタフェースは単純
//protocol server
type ProtocolServer interface {
ListenAndServerTLS() error
close() error
}websocket の実装も単純
func (srv *WSServer) ListenAndServeTLS() error{
return srv.server.ListenAndServeTLS("","")
}
func (srv *WSServer)Close() error{
if srv.server != nil{
return srv.server.Close()
}
return nil
}WSServer の核心は ServeHTTP。エッジの接入を扱う。流れ:

client インタフェースと websocket
client も単純。Connect はクラウドの server につなぎ Connection を返す。エッジはその connection で読み書きする。
//each protocol(websocket/quic) provides Connect
type Protocolclient interface{
Connect() (conn.Connection,error)
}websocket Client は dial でクラウドへつなぐ。成功したら Callback があれば呼び、connection を初期化して返す。
//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 は CloudCore のモジュール。エッジの接入と雲辺データ転送。Controllers と EdgeCore の間の仲介。下行メッセージ(中に k8s リソースイベント、pod update など)をエッジへ配り、エッジからの状態を対応する controllers へ渡す。位置:

内部の主なコードモジュール:

- HTTP server:エッジへ証明書の入口。CA、発行、ローテ
- WebSocket server:オンオフ可。エッジの WebSocket 接入
- QUIC server:オンオフ可。エッジの QUIC 接入
- CSI socket server:クラウドで csi driver と話す
- Token manager:エッジ接入トークン。既定 12h ローテ
- Certificate manager:エッジ証明書の発行とローテ
- message handler:エッジ接入とメッセージ配送
- node session manager:エッジセッションのライフサイクル
- message dispatcher:上り下りメッセージの配送
Cloudhub の起動
CloudCore 起動時に登録し、beehive が Start() を呼ぶ
cloudhub.Register(c.modules.Cloudhub)起動時、まず dispatcher.DispatchDownstream を回して下行を非同期配送。次に証明書。未設定なら CA とサーバ証明書を自動生成し、あとで WebSocket / Quic / HTTP の安全通信に使う。token manager を起動し、エッジ用トークンを作り自動ローテを始める。StartHTTPServer() は主に EdgeCore の証明書申請待ち。それから cloudhub 本体:viaduct でサーバを起動し、EdgeCore の接続を待つ。tcp の WebSocket か udp の QUIC。CSI が要れば CSI socket server も起動。
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
核心はエッジ接入とメッセージ配送。内部構成:

下行メッセージの送り方
エッジへの下行には二モード。配送と node session の処理に直結する:
ACK:エッジが下行を受け、ローカルストアに正しく保存したあと、クラウドへ ACK を返す。ACK が来なければ未処理とみなし、ACK までリトライする。
NO-ACK:エッジは ACK を返さない。クラウドは届いて処理されたとみなす。落ちることがある。同期メッセージへの応答に使うことが多く、エッジが応答を受け取らなければリトライする。
エッジノードの接入
主な論理は 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 が接入を扱う。WebSocket の例:viaduct で WS server が上がったあと、エッジが来ると ServeHTTP が http を websocket に upgrade し、Connection を初期化。HandleConnection は渡された connection で:
- 初期化前の検査。ノード数上限など。
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
}- nodeMessagePool を初期化し、MessageDispatcher のハッシュに入れ、下行の置き場にする。
//init node message pool and add to the dispatcher
nodeMessagePool := common.InitNodeMessagePool(nodeID)
mh.MessageDispatcher.AddNodeMessagePool(nodeID,nodeMessagePool)nodeMessagePool は下行キュー。エッジが接入するたびに一つ。上の二モードに対応し、ACK 用と 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
- nodeSession を初期化し、SessionManager のハッシュに入れ、起動する
//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()エッジノードごとに nodeSession。接続セッションの抽象。SessionManager はこの CloudHub につながる全セッションを持つ。起動するとそのノードに要る goroutine を全部起こす: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()
}上り下りの配送
CloudHub では上り下りは比較的単純。messageHandler の HandleMessage。下の viaduct が MessageContainer にパースし、message が入っている。HandleMessage は軽く検査し、MessageDispatcher.DispatchUpstream で edgeController や deviceController などへ。
//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})
}下行、ACK の場合:
- KubeEdge は K8s objectSync CRD で、Edge へ送り切ったリソースの最新 resourceVersion を持つ。Cloudhub 起動時、これから送る resourceVersion と送り済みを比べ、古いメッセージを送らない。
- EdgeController や devicecontroller などが Cloudhub へ送る。MessageDispatcher はノード名で対応する NodeMessagePool へ入れ、resource などから送りモードを選ぶ。キューに入れるとき objectSync CR を引き、送り済みの最新 resourceVersion と比べ、重複を避ける。
- そのノードの SendAckMessage が NodeMessagePool から順に取り、エッジへ送り、メッセージ ID を ACK channel に置く。エッジから ACK が来ると channel が通知し、resourceVersion を objectSync CR に保存し、次を送る。
- EdgeCore はまずローカルストアに保存し、ACK をクラウドへ返す。この間隔で ACK が来なければ最大 5 回再送。5 回とも失敗ならイベントを捨てる。
- CloudCore の SyncController が失敗イベントを扱う。エッジが受け取っていても、ACK が途中で落ちることがある。そのとき SyncController がもう一度 Cloudhub へ送り、下行をやり直す。成功するまで。
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
エッジ計算ではネットワークが不安定なことが多く、雲辺は切れやすい。データが落ちるリスクがある。synccontroller は CloudCore のモジュールで、メッセージの確実な送信を守る。KubeEdge では objectSync で雲辺のメッセージ状態を永続化する。同期の途中、クラウドは各エッジで成功した最新 ResourceVersion を記録し、CR として K8s に残す。これでクラウド障害やエッジのオフライン再起動のあとも順序と連続が保たれ、古いメッセージの再送で状態がずれない。同時に周期的に雲辺データを見て、一貫を保つ。各エッジの同期状態を K8s のリソースと比べ、ずれをエッジへ送り、最終一貫にする。
CloudCore 起動時に登録し、beehive の Start() で動く。
synccontroller.Register(c.Modules.syncController)起動すると 5 秒間隔の周期チェックが始まる。
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 は名前空間スコープのオブジェクトを持つ。名前はノード名とオブジェクト UUID。SyncController は保存した ObjectSync の送り済み resourceVersion を K8s 上のオブジェクトと定期比較し、リトライや削除を起こす。cloudhub が NodeMessagePool にイベントを足すとき、プール内の対応オブジェクトと比べる。プールの方が新しければ捨て、そうでなければエッジへ送る。

edgehub
EdgeHub は WebSocket または QUIC のクライアント。CloudCore とやり取りする。クラウドのリソース更新の同期、エッジホストとデバイス状態の報告など。
EdgeCore 起動時に beehive で登録し、初期化する。
func Register(eh *v1alpha2.EdgeHub,nodeName string){
config.InitConfigure(eh,nodeName)
core.Register(newEdgeHub(eh.Enable))
}起動コード:
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
起動の流れ:
- 証明書の初期化。cloudcore から申請(ローカルが正しく設定されていればそれを使う)。ローテを開始し、ループへ
- eh.initial() で eh.chClient を作り、eh.chClient.Init()。viaduct で websocket/quic の connection を張る
- eh.pubConnectInfo(true) で、つながったことを edgecore の各モジュールへ放送
- 三つの goroutine:
- routeToEdge
- routeToCloud
- keepalive
routeToEdge:クラウドから来たメッセージを受ける。同期応答なら beehive sendResp。そうでなければ、メッセージの group へ送る。
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:エッジ側の他 module からのメッセージを受け、websocket/quic client でクラウドへ送る
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:ハートビート周期でクラウドへ ping
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)
}
}- 雲辺の転送でエラーが出ると、エッジは websocket/quic client を再 init し、クラウドとつなぎ直す。