diff options
| -rw-r--r-- | cloud/pkg/cloudhub/cloudhub.go | 5 | ||||
| -rw-r--r-- | cloud/pkg/cloudstream/cloudstream.go | 18 | ||||
| -rw-r--r-- | cloud/pkg/cloudstream/tunnelserver.go | 22 | ||||
| -rw-r--r-- | edge/pkg/edgehub/edgehub.go | 4 | ||||
| -rw-r--r-- | edge/pkg/edgestream/edgestream.go | 42 | ||||
| -rw-r--r-- | pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go | 7 | ||||
| -rw-r--r-- | pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go | 16 |
7 files changed, 62 insertions, 52 deletions
diff --git a/cloud/pkg/cloudhub/cloudhub.go b/cloud/pkg/cloudhub/cloudhub.go index 2b99efd6f..3c6c13461 100644 --- a/cloud/pkg/cloudhub/cloudhub.go +++ b/cloud/pkg/cloudhub/cloudhub.go @@ -20,6 +20,8 @@ import ( "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1" ) +var DoneTLSTunnelCerts = make(chan bool, 1) + type cloudHub struct { enable bool } @@ -69,6 +71,9 @@ func (a *cloudHub) Start() { if err := httpserver.PrepareAllCerts(); err != nil { klog.Fatal(err) } + // TODO: Will improve in the future + DoneTLSTunnelCerts <- true + close(DoneTLSTunnelCerts) // generate Token if err := httpserver.GenerateToken(); err != nil { diff --git a/cloud/pkg/cloudstream/cloudstream.go b/cloud/pkg/cloudstream/cloudstream.go index 96b5d2098..bd8d52968 100644 --- a/cloud/pkg/cloudstream/cloudstream.go +++ b/cloud/pkg/cloudstream/cloudstream.go @@ -18,6 +18,7 @@ package cloudstream import ( "github.com/kubeedge/beehive/pkg/core" + "github.com/kubeedge/kubeedge/cloud/pkg/cloudhub" "github.com/kubeedge/kubeedge/cloud/pkg/cloudstream/config" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1" ) @@ -51,13 +52,18 @@ func (s *cloudStream) Group() string { } func (s *cloudStream) Start() { - ts := newTunnelServer() - // start new tunnel server - go ts.Start() + // TODO: Will improve in the future + ok := <-cloudhub.DoneTLSTunnelCerts + if ok { + ts := newTunnelServer() - server := newStreamServer(ts) - // start stream server to accepet kube-apiserver connection - go server.Start() + // start new tunnel server + go ts.Start() + + server := newStreamServer(ts) + // start stream server to accepet kube-apiserver connection + go server.Start() + } } func (s *cloudStream) Enable() bool { diff --git a/cloud/pkg/cloudstream/tunnelserver.go b/cloud/pkg/cloudstream/tunnelserver.go index 9a4bc3c59..e2a7b9a46 100644 --- a/cloud/pkg/cloudstream/tunnelserver.go +++ b/cloud/pkg/cloudstream/tunnelserver.go @@ -19,19 +19,20 @@ package cloudstream import ( "crypto/tls" "crypto/x509" + "encoding/pem" "fmt" "io/ioutil" "net/http" "sync" "time" - "github.com/kubeedge/kubeedge/pkg/stream" - "github.com/emicklei/go-restful" "github.com/gorilla/websocket" "k8s.io/klog" + hubconfig "github.com/kubeedge/kubeedge/cloud/pkg/cloudhub/config" "github.com/kubeedge/kubeedge/cloud/pkg/cloudstream/config" + "github.com/kubeedge/kubeedge/pkg/stream" ) type TunnelServer struct { @@ -102,24 +103,27 @@ func (s *TunnelServer) Start() { s.installDefaultHandler() data, err := ioutil.ReadFile(config.Config.TLSTunnelCAFile) if err != nil { - klog.Fatalf("Read tls tunnel ca file error %v", err) - return + data = hubconfig.Config.Ca + klog.Info("Succeeded in loading TLSTunnelCAFile from the secret") } pool := x509.NewCertPool() - pool.AppendCertsFromPEM(data) + pool.AppendCertsFromPEM(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: data})) cert, err := ioutil.ReadFile(config.Config.TLSTunnelCertFile) if err != nil { - klog.Fatalf("Read cert file %v error %v", config.Config.TLSTunnelCertFile, err) + cert = hubconfig.Config.Cert + klog.Info("Succeeded in loading TLSTunnelCertFile from the secret") } key, err := ioutil.ReadFile(config.Config.TLSTunnelPrivateKeyFile) if err != nil { - klog.Fatalf("Read key file %v error %v", config.Config.TLSTunnelPrivateKeyFile, err) + key = hubconfig.Config.Key + klog.Info("Succeeded in loading TLSTunnelPrivateKeyFile from the secret") } - certificate, err := tls.X509KeyPair(cert, key) + certificate, err := tls.X509KeyPair(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: cert}), pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: key})) if err != nil { panic(err) } + klog.Info("Succeeded in loading TLSTunnelCert and Key") tunnelServer := &http.Server{ Addr: fmt.Sprintf(":%d", config.Config.TunnelPort), @@ -129,7 +133,7 @@ func (s *TunnelServer) Start() { Certificates: []tls.Certificate{certificate}, ClientAuth: tls.RequireAndVerifyClientCert, MinVersion: tls.VersionTLS12, - CipherSuites: []uint16{tls.TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256}, + CipherSuites: []uint16{tls.TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256}, }, } klog.Infof("Prepare to start tunnel server ...") diff --git a/edge/pkg/edgehub/edgehub.go b/edge/pkg/edgehub/edgehub.go index d91b196b9..721c93184 100644 --- a/edge/pkg/edgehub/edgehub.go +++ b/edge/pkg/edgehub/edgehub.go @@ -22,6 +22,8 @@ const ( ModuleNameEdgeHub = "websocket" ) +var HasTLSTunnelCerts = make(chan bool, 1) + //EdgeHub defines edgehub object structure type EdgeHub struct { chClient clients.Adapter @@ -77,6 +79,8 @@ func (eh *EdgeHub) Start() { return } } + HasTLSTunnelCerts <- true + close(HasTLSTunnelCerts) for { select { diff --git a/edge/pkg/edgestream/edgestream.go b/edge/pkg/edgestream/edgestream.go index 0602de68e..8e5116669 100644 --- a/edge/pkg/edgestream/edgestream.go +++ b/edge/pkg/edgestream/edgestream.go @@ -27,6 +27,7 @@ import ( "github.com/kubeedge/beehive/pkg/core" beehiveContext "github.com/kubeedge/beehive/pkg/core/context" + "github.com/kubeedge/kubeedge/edge/pkg/edgehub" "github.com/kubeedge/kubeedge/edge/pkg/edgestream/config" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1" "github.com/kubeedge/kubeedge/pkg/stream" @@ -71,32 +72,33 @@ func (e *edgestream) Enable() bool { } func (e *edgestream) Start() { - serverURL := url.URL{ Scheme: "wss", Host: config.Config.TunnelServer, Path: "/v1/kubeedge/connect", } - - cert, err := tls.LoadX509KeyPair(config.Config.TLSTunnelCertFile, config.Config.TLSTunnelPrivateKeyFile) - if err != nil { - klog.Fatalf("Failed to load x509 key pair: %v", err) - } - - 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) + // TODO: Will improve in the future + ok := <-edgehub.HasTLSTunnelCerts + if ok { + cert, err := tls.LoadX509KeyPair(config.Config.TLSTunnelCertFile, config.Config.TLSTunnelPrivateKeyFile) if err != nil { - klog.Errorf("TLSClientConnect error %v", err) + klog.Fatalf("Failed to load x509 key pair: %v", err) + } + 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) + if err != nil { + klog.Errorf("TLSClientConnect error %v", err) + } } } } diff --git a/pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go b/pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go index 4d61b7639..69c84d55e 100644 --- a/pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go +++ b/pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go @@ -24,6 +24,7 @@ import ( "k8s.io/apimachinery/pkg/util/validation/field" componentbaseconfig "k8s.io/component-base/config" + "k8s.io/klog" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1" utilvalidation "github.com/kubeedge/kubeedge/pkg/util/validation" @@ -149,13 +150,13 @@ func ValidateModuleCloudStream(d v1alpha1.CloudStream) field.ErrorList { allErrs := field.ErrorList{} if !utilvalidation.FileIsExist(d.TLSTunnelPrivateKeyFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelPrivateKeyFile"), d.TLSTunnelPrivateKeyFile, "TLSTunnelPrivateKeyFile not exist")) + klog.Warningf("TLSTunnelPrivateKeyFile does not exist in %s, will load from secret", d.TLSTunnelPrivateKeyFile) } if !utilvalidation.FileIsExist(d.TLSTunnelCertFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelCertFile"), d.TLSTunnelCertFile, "TLSTunnelCertFile not exist")) + klog.Warningf("TLSTunnelCertFile does not exist in %s, will load from secret", d.TLSTunnelCertFile) } if !utilvalidation.FileIsExist(d.TLSTunnelCAFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelCAFile"), d.TLSTunnelCAFile, "TLSTunnelCAFile not exist")) + klog.Warningf("TLSTunnelCAFile does not exist in %s, will load from secret", d.TLSTunnelCAFile) } if !utilvalidation.FileIsExist(d.TLSStreamPrivateKeyFile) { diff --git a/pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go b/pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go index 48a9e3c50..25025bdce 100644 --- a/pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go +++ b/pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go @@ -152,21 +152,9 @@ func ValidateModuleEdgeMesh(m v1alpha1.EdgeMesh) field.ErrorList { // ValidateModuleEdgeStream validates `m` and returns an errorList if it is invalid func ValidateModuleEdgeStream(m v1alpha1.EdgeStream) field.ErrorList { - if !m.Enable { - return field.ErrorList{} - } allErrs := field.ErrorList{} - if !utilvalidation.FileIsExist(m.TLSTunnelCAFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelCAFile"), - m.TLSTunnelCAFile, "TLSTunnelCAFile file not exist")) - } - if !utilvalidation.FileIsExist(m.TLSTunnelCertFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelCertFile"), - m.TLSTunnelCertFile, "TLSTunnelCertFile file not exist")) - } - if !utilvalidation.FileIsExist(m.TLSTunnelPrivateKeyFile) { - allErrs = append(allErrs, field.Invalid(field.NewPath("TLSTunnelPrivateKeyFile"), - m.TLSTunnelPrivateKeyFile, "TLSTunnelPrivateKeyFile file not exist")) + if !m.Enable { + return allErrs } return allErrs } |
