作者/古强
Stream
Stream是KubeEdge中提供云边隧道的模块,目前支持ApiServer向Kubelet发起的containerLog、exec和metrics请求。云边隧道基于WebSocket建造,支持双向传输和流式传输。
架构
整个Stream功能由CloudStream和EdgeStream两部分组成,从名字便可以看出它们各自运行在何处。

除CloudStream和EdgeStream模块,Stream功能的正常使用要配置iptables规则。所有向边缘节点DaemonEndpoint的请求都会被转发到CloudStream的StreamPort,经过内部处理后从边缘节点与TunnelPort建立的云边隧道到达EdgeStream,由EdgeStream去向Edged发起请求。
概念
Session:边缘节点和TunnelServer间长连接的云端抽象,TunnelServer包含若干Session。
APIServerConnection:单次API Server请求的云端抽象,一个Session包含若干个正在处理的APIServerConnection。
TunnelSession:边缘节点和TunnelServer间长连接的边缘侧抽象,每个EdgeStream模块只有一个。
EdgedConnection:单次API Server请求的边缘侧抽象,一个TunnelSession包含若干正在处理的Edged Connection。
Message:云边实际传输的数据封包,由连接ID、数据类型和数据三部分组成。
每个云边隧道由一个Session和一个TunnelSession配对组成,提供网络传输,上层只需收发Message数据包而无需关注传输细节。
CloudStream
CloudStream启动时在协程中同时运行TunnelServer和StreamServer。
func (s *cloudStream) Start() {
// TODO: Will improve in the future
ok := <-cloudhub.DoneTLSTunnelCerts
if ok {
ts := newTunnelServer()
// start new tunnel server
go ts.Start()
server := newStreamServer(ts)
// start stream server to accepet kube-apiserver connection
go server.Start()
}
我们首先看一下TunnelServer的Start方法:
func (s *TunnelServer) Start() {
s.installDefaultHandler() // 设置路由
var data []byte
var key []byte
var cert []byte
... // 根据配置读取data
pool := x509.NewCertPool()
pool.AppendCertsFromPEM(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: data}))
... // 根据配置读取key和cert
certificate, err := tls.X509KeyPair(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: cert}), pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: key}))
...
tunnelServer := &http.Server{
Addr: fmt.Sprintf(":%d", config.Config.TunnelPort),
Handler: s.container,
TLSConfig: &tls.Config{
ClientCAs: pool,
Certificates: []tls.Certificate{certificate},
ClientAuth: tls.RequireAndVerifyClientCert,
MinVersion: tls.VersionTLS12,
CipherSuites: []uint16{tls.TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256},
},
}
klog.Infof("Prepare to start tunnel server ...")
err = tunnelServer.ListenAndServeTLS("", "")
if err != nil {
klog.Fatalf("Start tunnelServer error %v\n", err)
return
}
}TunnelServer在启动时,首先通过installDefaultHandler方法注册路由/v1/kubeedge/connect,之后根据配置设置证书并启动Server监听。
当请求来到/v1/kubeedge/connect后,会交给 TunnelServer 的connect方法处理:
func (s *TunnelServer) connect(r *restful.Request, w *restful.Response) {
hostNameOverride := r.HeaderParameter(stream.SessionKeyHostNameOveride)
interalIP := r.HeaderParameter(stream.SessionKeyInternalIP)
if interalIP == "" {
interalIP = strings.Split(r.Request.RemoteAddr, ":")[0]
}
con, err := s.upgrader.Upgrade(w, r.Request, nil)
if err != nil {
return
}
klog.Infof("get a new tunnel agent hostname %v, internalIP %v", hostNameOverride, interalIP)
session := &Session{
tunnel: stream.NewDefaultTunnel(con),
apiServerConn: make(map[uint64]APIServerConnection),
apiConnlock: &sync.Mutex{},
sessionID: hostNameOverride,
}
s.addSession(hostNameOverride, session)
s.addSession(interalIP, session)
session.Serve()
}Connect会从请求中提取出客户端的主机名和IP地址,之后将HTTP请求升级为WebSocket请求,最后将WebSocket和主机信息包装成Session并缓存下来,最后调用Serve方法。
接下来我们看一下StreamServer的Start方法:
func (s *StreamServer) Start() {
s.installDebugHandler() // 设置路由
pool := x509.NewCertPool()
data, err := ioutil.ReadFile(config.Config.TLSStreamCAFile)
...
pool.AppendCertsFromPEM(data)
streamServer := &http.Server{
Addr: fmt.Sprintf(":%d", config.Config.StreamPort),
Handler: s.container,
TLSConfig: &tls.Config{
ClientCAs: pool,
// Populate PeerCertificates in requests, but don't reject connections without verified certificates
ClientAuth: tls.RequestClientCert,
},
}
klog.Infof("Prepare to start stream server ...")
err = streamServer.ListenAndServeTLS(config.Config.TLSStreamCertFile, config.Config.TLSStreamPrivateKeyFile)
...
}StreamServer在启动时,首先通过installDebugHandler方法注册了多个路由和对应的handler,例如路由/containerLogs/{podNamespace}/{podID}/{containerName}。然后便是配置证书并启动server监听请求。
当 API Server 的请求来到 StreamServer 后,根据不同的路由匹配相应handler处理,这里我们以容器日志/containerLogs/{podNamespace}/{podID}/{containerName}为例,handler为getContainerLogs。
func (s *StreamServer) getContainerLogs(r *restful.Request, w *restful.Response) {
var err error
defer func() {
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
klog.Errorf(err.Error())
}
}()
sessionKey := strings.Split(r.Request.Host, ":")[0]
session, ok := s.tunnel.getSession(sessionKey)
...
w.Header().Set("Transfer-Encoding", "chunked")
w.WriteHeader(http.StatusOK)
if _, ok := w.ResponseWriter.(http.Flusher); !ok {
...
return
}
fw := flushwriter.Wrap(w.ResponseWriter)
logConnection, err := session.AddAPIServerConnection(s, &ContainerLogsConnection{
r: r,
flush: fw,
session: session,
ctx: r.Request.Context(),
edgePeerStop: make(chan struct{}),
})
...
defer func() {
session.DeleteAPIServerConnection(logConnection)
klog.Infof("Delete %s from %s", logConnection.String(), session.String())
}()
if err := logConnection.Serve(); err != nil {
...
return
}
}首先从请求中取出主机地址,从TunnelServer中寻找相应的Session缓存,然后将本次请求包装为一次APIServerConnection,最后调用APIServerConnection.Serve()方法。
APIServerConnection是一个接口,接下来我们看一下在获取容器日志时,Serve()方法是如何实现的。
func (l *ContainerLogsConnection) Serve() error {
defer func() {
klog.Infof("%s end successful", l.String())
}()
// first send connect message
if _, err := l.SendConnection(); err != nil {
...
return err
}
for {
select {
case <-l.ctx.Done():
// if apiserver request end, send close message to edge
msg := stream.NewMessage(l.MessageID, stream.MessageTypeRemoveConnect, nil)
for retry := 0; retry < 3; retry++ {
if err := l.WriteToTunnel(msg); err != nil {
...
} else {
break
}
}
klog.Infof("%s send close message to edge successfully", l.String())
return nil
case <-l.EdgePeerDone():
klog.Infof("%s find edge peer done, so stop this connection", l.String())
return nil
}
}
}在这个Serve方法中,通过SendConnection向EdgeStream发送了一条消息,之后便是针对请求结束和连接断开的处理,没有接收返回数据的逻辑。
那么,让我们回头审视一下Session.Serve()方法:
// Serve read tunnel message ,and write to specific apiserver connection
func (s *Session) Serve() {
defer s.Close()
for {
t, r, err := s.tunnel.NextReader()
...
message, err := stream.ReadMessageFromTunnel(r)
...
if err := s.ProxyTunnelMessageToApiserver(message); err != nil {
...
continue
}
}
}在这个方法中,Session会不断地读取WebSocket消息,并调用ProxyTunnelMessageToApiserver进行处理,让我们继续看看ProxyTunnelMessageToApiserver做了什么。
func (s *Session) ProxyTunnelMessageToApiserver(message *stream.Message) error {
kubeCon, ok := s.apiServerConn[message.ConnectID]
...
switch message.MessageType {
case stream.MessageTypeRemoveConnect:
kubeCon.SetEdgePeerDone()
case stream.MessageTypeData:
for i := 0; i < len(message.Data); {
n, err := kubeCon.WriteToAPIServer(message.Data[i:]) // 带缓冲地写,可能需要多次才能完全发送消息
if err != nil {
return err
}
i += n
}
default:
}
return nil
}ProxyTunnelMessageToApiserver根据ConnectID获取对应的APIServerConnection接口实例,之后针对不同消息类型,调用不同的接口方法实现回调。对API Server请求响应也是在这里发出的。
现在,还剩下一环仍然缺失:EdgeStream如何接收到查询日志请求的?
让我们把目光放回SendConnection方法:
func (l *ContainerLogsConnection) SendConnection() (stream.EdgedConnection, error) {
connector := &stream.EdgedLogsConnection{
MessID: l.MessageID,
URL: *l.r.Request.URL,
Header: l.r.Request.Header,
}
connector.URL.Scheme = httpScheme
connector.URL.Host = net.JoinHostPort(defaultServerHost, fmt.Sprintf("%v", constants.ServerPort))
m, err := connector.CreateConnectMessage()
if err != nil {
return nil, err
}
if err := l.WriteToTunnel(m); err != nil {
klog.Errorf("%s write %s error %v", l.String(), connector.String(), err)
return nil, err
}
return connector, nil
}SendConnection把 API Server 请求的信息(保存在ContainerLogsConnection.r中)封装到EdgedLogsConnection中,再调用EdgedLogsConnection.CreateConnectMessage方法创建了一条消息并发送。
消息创建的具体过程就不展开了,我们只要知道在创建这条消息时指定了特殊的消息类别来代表查询日志的请求,并且查询的URL已经保存在Connector对象中就足够了。
云端TunnelServer中各对象的组合关系如下图:

EdgeStream
EdgeStream启动时会向TunnelServer注册自己。
func (e *edgestream) Start() {
serverURL := url.URL{
Scheme: "wss",
Host: config.Config.TunnelServer,
Path: "/v1/kubeedge/connect",
}
// TODO: Will improve in the future
ok := <-edgehub.HasTLSTunnelCerts
if ok
cert, err := tls.LoadX509KeyPair(config.Config.TLSTunnelCertFile, config.Config.TLSTunnelPrivateKeyFile)
...
tlsConfig := &tls.Config{
InsecureSkipVerify: true,
Certificates: []tls.Certificate{cert},
}
for range time.NewTicker(time.Second * 2).C {
select {
case <-beehiveContext.Done():
return
default:
}
err := e.TLSClientConnect(serverURL, tlsConfig)
...
}
}
}正常情况下,TLSClientConnect会一直监听WebSocket请求。一旦网络出现故障,外层的for循环便会在两秒后重试连接,设计比较巧妙。
TLSClientConnect方法的作用是根据配置构造WebSocket请求,并持续监听这个请求。
func (e *edgestream) TLSClientConnect(url url.URL, tlsConfig *tls.Config) error {
klog.Info("Start a new tunnel stream connection ...")
dial := websocket.Dialer{
TLSClientConfig: tlsConfig,
HandshakeTimeout: time.Duration(config.Config.HandshakeTimeout) * time.Second,
}
header := http.Header{}
header.Add(stream.SessionKeyHostNameOveride, e.hostnameOveride)
header.Add(stream.SessionKeyInternalIP, e.nodeIP)
con, _, err := dial.Dial(url.String(), header)
...
session := NewTunnelSession(con)
return session.Serve()
}接下来,我们看看Serve方法的逻辑。
func (s *TunnelSession) Serve() error {
ctx, cancel := context.WithCancel(context.Background())
defer func() {
cancel()
s.Close()
klog.Info("Close tunnel session successfully")
}()
go s.startPing(ctx)
for {
_, r, err := s.Tunnel.NextReader()
if err != nil {
klog.Errorf("Read Message error %v", err)
return err
}
mess, err := stream.ReadMessageFromTunnel(r)
if err != nil {
klog.Errorf("Get tunnel Message error %v", err)
return err
}
if mess.MessageType < stream.MessageTypeData {
go s.ServeConnection(mess)
}
s.WriteToLocalConnection(mess)
}
}Serve方法首先启动一个协程,不断向 TunnelServer 发送心跳消息,然后持续从 Tunnel 中读取 WebSocket 消息。如果读取到的消息类型枚举值小于数据类型(MessageTypeData=4),则新启动一个线程运行ServeConnection方法。无论是控制消息还是数据消息,都会由WriteToLocalConnection处理。
我们还是以日志查询为例,当查询日志的消息来到后,由于消息类型枚举值为1,会先在协程启动ServeConnection,最终执行的方法为serveLogsConnection。
func (s *TunnelSession) serveLogsConnection(m *stream.Message) error {
logCon := &stream.EdgedLogsConnection{
ReadChan: make(chan *stream.Message, 128),
}
if err := json.Unmarshal(m.Data, logCon); err != nil {
...
return err
}
s.AddLocalConnection(m.ConnectID, logCon)
return logCon.Serve(s.Tunnel)
}ServeLogsConnection方法创建一个EdgedLogsConnection对象,并从 WebSocket 消息中读取数据填充到EdgedLogsConnection,之后调用Serve方法。
func (l *EdgedLogsConnection) Serve(tunnel SafeWriteTunneler) error {
//connect edged
client := http.Client{}
req, err := http.NewRequest("GET", l.URL.String(), nil)
...
req.Header = l.Header
resp, err := client.Do(req)
...
defer resp.Body.Close()
reader := bufio.NewReader(resp.Body)
stop := make(chan struct{})
go func() {
for mess := range l.ReadChan {
if mess.MessageType == MessageTypeRemoveConnect {
klog.Infof("receive remove client id %v", mess.ConnectID)
close(stop)
return
}
}
}()
defer func() {
for retry := 0; retry < 3; retry++ {
msg := NewMessage(l.MessID, MessageTypeRemoveConnect, nil)
if err := msg.WriteTo(tunnel); err != nil {
...
} else {
break
}
}
}()
for {
select {
case <-stop:
klog.Infof("receive stop single, so stop logs scan ...")
return nil
default:
}
data := make([]byte, 256)
n, err := reader.Read(data)
if err != nil {
if err != io.EOF {
...
}
break
}
if n <= 0 {
continue
}
msg := NewMessage(l.MessID, MessageTypeData, data[:n])
err = msg.WriteTo(tunnel)
...
}
return nil
}查询容器日志不涉及双向通信,所以这里的代码显得有些冗余,我们来详细分析一下。
首先,向Edged请求容器日志,这里不是长连接,一次性拿到返回数据。
然后设置一个协程监听连接断开消息,在连接断开时退出读取HTTP返回数据的无限循环,这个消息将在API Server完成数据接收后由CloudStream发送过来。
之后,设置一个defer函数发送连接断开消息,该消息将发送给CloudStream,消息最多重试3次。
最后,在一个无限循环中通过分片把读取到的HTTP返回数据发送给CloudStream。
边缘侧EdgeStream中各对象的组合关系如下图:

总结
Stream模块目前只实现了logs、exec和metrics三类连接的隧道,且logs不支持follow属性。StreamServer监听的路径目前是硬编码不可配置的,需要扩展必须修改源码,不是很优雅。希望在后续的版本中,能够看到优化。