summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--cloud/pkg/cloudhub/cloudhub.go5
-rw-r--r--cloud/pkg/cloudstream/cloudstream.go18
-rw-r--r--cloud/pkg/cloudstream/tunnelserver.go22
-rw-r--r--edge/pkg/edgehub/edgehub.go4
-rw-r--r--edge/pkg/edgestream/edgestream.go42
-rw-r--r--pkg/apis/componentconfig/cloudcore/v1alpha1/validation/validation.go7
-rw-r--r--pkg/apis/componentconfig/edgecore/v1alpha1/validation/validation.go16
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
}