summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--cloud/pkg/common/client/client.go16
-rw-r--r--cloud/pkg/dynamiccontroller/application/application.go139
-rw-r--r--cloud/pkg/dynamiccontroller/application/application_test.go4
-rw-r--r--cloud/pkg/dynamiccontroller/application/eventhandler.go2
-rw-r--r--common/types/http.go4
-rw-r--r--edge/pkg/common/modules/names.go2
-rw-r--r--edge/pkg/edged/edged.go5
-rw-r--r--edge/pkg/edged/status/status_manager_test.go8
-rw-r--r--edge/pkg/edgehub/edgehub.go23
-rw-r--r--edge/pkg/edgestream/edgestream.go2
-rw-r--r--edge/pkg/metamanager/client/configmap.go3
-rw-r--r--edge/pkg/metamanager/client/metaclient.go6
-rw-r--r--edge/pkg/metamanager/client/node.go3
-rw-r--r--edge/pkg/metamanager/client/persistentvolume.go3
-rw-r--r--edge/pkg/metamanager/client/persistentvolumeclaim.go3
-rw-r--r--edge/pkg/metamanager/client/secret.go3
-rw-r--r--edge/pkg/metamanager/client/serviceaccount.go3
-rw-r--r--edge/pkg/metamanager/client/volumeattachment.go3
-rw-r--r--edge/pkg/metamanager/metamanager.go7
-rw-r--r--edge/pkg/metamanager/metamanager_test.go8
-rw-r--r--edge/pkg/metamanager/metaserver/config/config.go9
-rw-r--r--edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook/hooks.go1
-rw-r--r--edge/pkg/metamanager/metaserver/kubernetes/storage/storage.go67
-rw-r--r--edge/pkg/metamanager/metaserver/server.go102
-rw-r--r--edge/pkg/metamanager/process.go6
-rw-r--r--edge/pkg/metamanager/process_test.go6
-rw-r--r--edge/test/integration/metaserver/metaserver_suite_test.go4
-rw-r--r--edge/test/integration/metaserver/metaserver_test.go11
-rw-r--r--edge/test/test.go6
-rw-r--r--pkg/apis/componentconfig/edgecore/v1alpha1/default.go8
-rw-r--r--pkg/apis/componentconfig/edgecore/v1alpha1/types.go11
-rw-r--r--pkg/apis/componentconfig/meta/v1alpha1/types.go2
32 files changed, 361 insertions, 119 deletions
diff --git a/cloud/pkg/common/client/client.go b/cloud/pkg/common/client/client.go
index 5eb582733..5eb41e8da 100644
--- a/cloud/pkg/common/client/client.go
+++ b/cloud/pkg/common/client/client.go
@@ -36,6 +36,9 @@ var (
kubeClient kubernetes.Interface
crdClient crdClientset.Interface
dynamicClient dynamic.Interface
+ // authKubeConfig only contains master address and CA cert when init, it is used for
+ // generating a temporary kubeclient and validating user token once receive an application message.
+ authKubeConfig *rest.Config
)
func InitKubeEdgeClient(config *cloudcoreConfig.KubeAPIConfig) {
@@ -56,6 +59,15 @@ func InitKubeEdgeClient(config *cloudcoreConfig.KubeAPIConfig) {
crdKubeConfig := rest.CopyConfig(kubeConfig)
crdKubeConfig.ContentType = runtime.ContentTypeJSON
crdClient = crdClientset.NewForConfigOrDie(crdKubeConfig)
+
+ authKubeConfig, err = clientcmd.BuildConfigFromFlags(kubeConfig.Host, "")
+ if err != nil {
+ klog.Errorf("Failed to build config, err: %v", err)
+ os.Exit(1)
+ }
+ authKubeConfig.CAData = kubeConfig.CAData
+ authKubeConfig.CAFile = kubeConfig.CAFile
+ authKubeConfig.ContentType = runtime.ContentTypeJSON
})
}
@@ -70,3 +82,7 @@ func GetCRDClient() crdClientset.Interface {
func GetDynamicClient() dynamic.Interface {
return dynamicClient
}
+
+func GetAuthConfig() *rest.Config {
+ return authKubeConfig
+}
diff --git a/cloud/pkg/dynamiccontroller/application/application.go b/cloud/pkg/dynamiccontroller/application/application.go
index 399031433..beda38a46 100644
--- a/cloud/pkg/dynamiccontroller/application/application.go
+++ b/cloud/pkg/dynamiccontroller/application/application.go
@@ -6,9 +6,11 @@ import (
"encoding/json"
"errors"
"fmt"
+ "strings"
"sync"
"time"
+ authorizationv1 "k8s.io/api/authorization/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metainternalversion "k8s.io/apimachinery/pkg/apis/meta/internalversion"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -21,6 +23,8 @@ import (
apirequest "k8s.io/apiserver/pkg/endpoints/request"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/dynamic/dynamicinformer"
+ authorizationv1client "k8s.io/client-go/kubernetes/typed/authorization/v1"
+ "k8s.io/client-go/rest"
"k8s.io/klog/v2"
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
@@ -29,6 +33,7 @@ import (
"github.com/kubeedge/kubeedge/cloud/pkg/common/messagelayer"
"github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
"github.com/kubeedge/kubeedge/cloud/pkg/dynamiccontroller/filter"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
edgemodule "github.com/kubeedge/kubeedge/edge/pkg/common/modules"
metaserverconfig "github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/config"
@@ -96,16 +101,16 @@ type Application struct {
RespBody []byte
Subresource string
Error apierrors.StatusError
-
- ctx context.Context // to end app.Wait
- cancel context.CancelFunc
+ Token string
+ ctx context.Context // to end app.Wait
+ cancel context.CancelFunc
count uint64 // count the number of current citations
countLock sync.Mutex
timestamp time.Time // record the last closing time of application, only make sense when count == 0
}
-func newApplication(ctx context.Context, key string, verb applicationVerb, nodename, subresource string, option interface{}, reqBody interface{}) *Application {
+func newApplication(ctx context.Context, key string, verb applicationVerb, nodename, subresource string, option interface{}, reqBody interface{}) (*Application, error) {
var v1 metav1.ListOptions
if internal, ok := option.(metainternalversion.ListOptions); ok {
err := metainternalversion.Convert_internalversion_ListOptions_To_v1_ListOptions(&internal, &v1, nil)
@@ -115,6 +120,11 @@ func newApplication(ctx context.Context, key string, verb applicationVerb, noden
}
option = v1
}
+ token, ok := ctx.Value(commontypes.AuthorizationKey).(string)
+ if !ok {
+ klog.Errorf("unsupported Token type :%T", ctx.Value(commontypes.AuthorizationKey))
+ return nil, fmt.Errorf("unsupported Token type :%T", ctx.Value(commontypes.AuthorizationKey))
+ }
ctx2, cancel := context.WithCancel(ctx)
app := &Application{
Key: key,
@@ -124,6 +134,7 @@ func newApplication(ctx context.Context, key string, verb applicationVerb, noden
Status: PreApplying,
Option: toBytes(option),
ReqBody: toBytes(reqBody),
+ Token: token,
ctx: ctx2,
cancel: cancel,
count: 0,
@@ -131,7 +142,7 @@ func newApplication(ctx context.Context, key string, verb applicationVerb, noden
timestamp: time.Time{},
}
app.add()
- return app
+ return app, nil
}
func (a *Application) Identifier() string {
@@ -281,27 +292,29 @@ func NewApplicationAgent() *Agent {
return defaultAgent
}
-func (a *Agent) Generate(ctx context.Context, verb applicationVerb, option interface{}, obj runtime.Object) *Application {
+func (a *Agent) Generate(ctx context.Context, verb applicationVerb, option interface{}, obj runtime.Object) (*Application, error) {
key, err := metaserver.KeyFuncReq(ctx, "")
if err != nil {
klog.Errorf("%v", err)
- return &Application{}
+ return nil, err
}
info, ok := apirequest.RequestInfoFrom(ctx)
if !ok || !info.IsResourceRequest {
klog.Errorf("no request info in context")
- return &Application{}
+ return nil, fmt.Errorf("no request info in context")
+ }
+ app, err := newApplication(ctx, key, verb, a.nodeName, info.Subresource, option, obj)
+ if err != nil {
+ return nil, err
}
-
- app := newApplication(ctx, key, verb, a.nodeName, info.Subresource, option, obj)
store, ok := a.Applications.LoadOrStore(app.Identifier(), app)
if ok {
app = store.(*Application)
app.add()
- return app
+ return app, nil
}
- return app
+ return app, nil
}
func (a *Agent) Apply(app *Application) error {
@@ -337,7 +350,6 @@ func (a *Agent) Apply(app *Application) error {
func (a *Agent) doApply(app *Application) {
defer app.Call()
-
// encapsulate as a message
app.Status = InApplying
msg := model.NewMessage("").SetRoute(MetaServerSource, modules.DynamicControllerModuleGroup).FillBody(app)
@@ -348,7 +360,6 @@ func (a *Agent) doApply(app *Application) {
app.Reason = fmt.Sprintf("failed to access cloud Application center: %v", err)
return
}
-
retApp, err := msgToApplication(resp)
if err != nil {
app.Status = Failed
@@ -378,13 +389,13 @@ type Center struct {
Applications sync.Map
HandlerCenter
messageLayer messagelayer.MessageLayer
- kubeclient dynamic.Interface
+ authConfig *rest.Config
}
func NewApplicationCenter(dynamicSharedInformerFactory dynamicinformer.DynamicSharedInformerFactory) *Center {
a := &Center{
HandlerCenter: NewHandlerCenter(dynamicSharedInformerFactory),
- kubeclient: client.GetDynamicClient(),
+ authConfig: client.GetAuthConfig(),
messageLayer: messagelayer.DynamicControllerMessageLayer(),
}
return a
@@ -429,6 +440,7 @@ func (c *Center) Process(msg model.Message) {
klog.Errorf("failed to translate msg to Application: %v", err)
return
}
+
klog.Infof("[metaserver/ApplicationCenter] get a Application %v", app.String())
resp, err := c.ProcessApplication(app)
@@ -441,6 +453,66 @@ func (c *Center) Process(msg model.Message) {
klog.Infof("[metaserver/applicationCenter]successfully to process Application(%+v)", app)
}
+func (c *Center) generateNewConfig(raw string) (*rest.Config, error) {
+ parts := strings.SplitN(raw, " ", 3)
+ if len(parts) < 2 || strings.ToLower(parts[0]) != "bearer" || len(parts[1]) <= 0 {
+ return nil, fmt.Errorf("invalid request token format or length: %v", len(parts))
+ }
+ authConfig := rest.CopyConfig(c.authConfig)
+ authConfig.BearerToken = parts[1]
+ return authConfig, nil
+}
+
+func (c *Center) createAuthClient(app *Application) (authorizationv1client.AuthorizationV1Interface, error) {
+ authConfig, err := c.generateNewConfig(app.Token)
+ if err != nil {
+ return nil, err
+ }
+ return authorizationv1client.NewForConfigOrDie(authConfig), nil
+}
+
+func (c *Center) createKubeClient(app *Application) (dynamic.Interface, error) {
+ authConfig, err := c.generateNewConfig(app.Token)
+ if err != nil {
+ return nil, err
+ }
+ return dynamic.NewForConfigOrDie(authConfig), nil
+}
+
+func (c *Center) authorizeApplication(app *Application, gvr schema.GroupVersionResource, namespace string, name string) error {
+ tmpAuthClient, err := c.createAuthClient(app)
+ if err != nil {
+ return err
+ }
+ sar := &authorizationv1.SelfSubjectAccessReview{
+ Spec: authorizationv1.SelfSubjectAccessReviewSpec{
+ ResourceAttributes: &authorizationv1.ResourceAttributes{
+ Namespace: namespace,
+ Verb: string(app.Verb),
+ Group: gvr.Group,
+ Resource: gvr.Resource,
+ Name: name,
+ Subresource: app.Subresource,
+ },
+ },
+ }
+ response, err := tmpAuthClient.SelfSubjectAccessReviews().Create(context.TODO(), sar, metav1.CreateOptions{})
+ if err != nil {
+ return err
+ }
+ if response.Status.Allowed {
+ return nil
+ }
+ var errMsg = fmt.Sprintf("resource %v authorize failed.", gvr)
+ if len(response.Status.Reason) > 0 {
+ errMsg += fmt.Sprintf("reason: %v.", response.Status.Reason)
+ }
+ if len(response.Status.EvaluationError) > 0 {
+ errMsg += fmt.Sprintf("evaluation error: %v.", response.Status.EvaluationError)
+ }
+ return fmt.Errorf(errMsg)
+}
+
// ProcessApplication processes application by re-translating it to kube-api request with kube client,
// which will be processed and responded by apiserver eventually.
// Specially if app.verb == watch, it transforms app to a listener and register it to HandlerCenter, rather
@@ -448,20 +520,33 @@ func (c *Center) Process(msg model.Message) {
// push them to edge node.
func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
app.Status = InProcessing
-
gvr, ns, name := metaserver.ParseKey(app.Key)
+ var kubeClient dynamic.Interface
+ var err error
+ if app.Verb != Watch {
+ kubeClient, err = c.createKubeClient(app)
+ if err != nil {
+ klog.Errorf("create kube client error: %v", err)
+ return nil, err
+ }
+ }
switch app.Verb {
case List:
var option = new(metav1.ListOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
- list, err := c.kubeclient.Resource(app.GVR()).Namespace(app.Namespace()).List(context.TODO(), *option)
+ list, err := kubeClient.Resource(gvr).Namespace(ns).List(context.TODO(), *option)
if err != nil {
- return nil, fmt.Errorf("successfully to add listener but failed to get current list, %v", err)
+ return nil, fmt.Errorf("get current list error: %v", err)
}
return list, nil
case Watch:
+ err := c.authorizeApplication(app, gvr, ns, name)
+ if err != nil {
+ klog.Errorf("authorize application error: %v", err)
+ return nil, err
+ }
var option = new(metav1.ListOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
@@ -475,7 +560,7 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
if err := app.OptionTo(option); err != nil {
return nil, err
}
- retObj, err := c.kubeclient.Resource(gvr).Namespace(ns).Get(context.TODO(), name, *option)
+ retObj, err := kubeClient.Resource(gvr).Namespace(ns).Get(context.TODO(), name, *option)
if err != nil {
return nil, err
}
@@ -492,9 +577,9 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
var retObj interface{}
var err error
if app.Subresource == "" {
- retObj, err = c.kubeclient.Resource(gvr).Namespace(ns).Create(context.TODO(), obj, *option)
+ retObj, err = kubeClient.Resource(gvr).Namespace(ns).Create(context.TODO(), obj, *option)
} else {
- retObj, err = c.kubeclient.Resource(gvr).Namespace(ns).Create(context.TODO(), obj, *option, app.Subresource)
+ retObj, err = kubeClient.Resource(gvr).Namespace(ns).Create(context.TODO(), obj, *option, app.Subresource)
}
if err != nil {
return nil, err
@@ -505,7 +590,7 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
if err := app.OptionTo(&option); err != nil {
return nil, err
}
- if err := c.kubeclient.Resource(gvr).Namespace(ns).Delete(context.TODO(), name, *option); err != nil {
+ if err := kubeClient.Resource(gvr).Namespace(ns).Delete(context.TODO(), name, *option); err != nil {
return nil, err
}
return nil, nil
@@ -521,9 +606,9 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
var retObj interface{}
var err error
if app.Subresource == "" {
- retObj, err = c.kubeclient.Resource(gvr).Namespace(ns).Update(context.TODO(), obj, *option)
+ retObj, err = kubeClient.Resource(gvr).Namespace(ns).Update(context.TODO(), obj, *option)
} else {
- retObj, err = c.kubeclient.Resource(gvr).Namespace(ns).Update(context.TODO(), obj, *option, app.Subresource)
+ retObj, err = kubeClient.Resource(gvr).Namespace(ns).Update(context.TODO(), obj, *option, app.Subresource)
}
if err != nil {
return nil, err
@@ -538,7 +623,7 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
if err := app.ReqBodyTo(obj); err != nil {
return nil, err
}
- retObj, err := c.kubeclient.Resource(gvr).Namespace(ns).UpdateStatus(context.TODO(), obj, *option)
+ retObj, err := kubeClient.Resource(gvr).Namespace(ns).UpdateStatus(context.TODO(), obj, *option)
if err != nil {
return nil, err
}
@@ -548,7 +633,7 @@ func (c *Center) ProcessApplication(app *Application) (interface{}, error) {
if err := app.OptionTo(pi); err != nil {
return nil, err
}
- retObj, err := c.kubeclient.Resource(gvr).Namespace(ns).Patch(context.TODO(), pi.Name, pi.PatchType, pi.Data, pi.Options, pi.Subresources...)
+ retObj, err := kubeClient.Resource(gvr).Namespace(ns).Patch(context.TODO(), pi.Name, pi.PatchType, pi.Data, pi.Options, pi.Subresources...)
if err != nil {
return nil, err
}
diff --git a/cloud/pkg/dynamiccontroller/application/application_test.go b/cloud/pkg/dynamiccontroller/application/application_test.go
index 60aaff9a0..7aa370991 100644
--- a/cloud/pkg/dynamiccontroller/application/application_test.go
+++ b/cloud/pkg/dynamiccontroller/application/application_test.go
@@ -8,6 +8,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
apirequest "k8s.io/apiserver/pkg/endpoints/request"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
metaserverconfig "github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/config"
)
@@ -32,8 +33,9 @@ func TestApplicationGC(t *testing.T) {
Resource: "nodes",
}
ctx := apirequest.WithRequestInfo(context.Background(), requestInfo)
+ ctx = context.WithValue(ctx, commontypes.AuthorizationKey, "Bearer xxxx")
- app := a.Generate(ctx, "get", metav1.GetOptions{}, nil)
+ app, _ := a.Generate(ctx, "get", metav1.GetOptions{}, nil)
app.countLock.Lock()
app.count = 0
app.countLock.Unlock()
diff --git a/cloud/pkg/dynamiccontroller/application/eventhandler.go b/cloud/pkg/dynamiccontroller/application/eventhandler.go
index f56af37fe..26c5953fb 100644
--- a/cloud/pkg/dynamiccontroller/application/eventhandler.go
+++ b/cloud/pkg/dynamiccontroller/application/eventhandler.go
@@ -70,7 +70,7 @@ func (c *handlerCenter) ForResource(gvr schema.GroupVersionResource) *CommonReso
return handler
}
-// dispatch listeners to corresponding CommonResourceEventHandler according it's gvr
+// AddListener dispatch listeners to corresponding CommonResourceEventHandler according it's gvr
func (c *handlerCenter) AddListener(s *SelectorListener) error {
return c.ForResource(s.gvr).AddListener(s)
}
diff --git a/common/types/http.go b/common/types/http.go
index ba704ca6b..21a4ed7c2 100644
--- a/common/types/http.go
+++ b/common/types/http.go
@@ -16,3 +16,7 @@ type HTTPResponse struct {
StatusCode int `json:"status_code"`
Body []byte `json:"body"`
}
+
+const (
+ AuthorizationKey = "Authorization"
+)
diff --git a/edge/pkg/common/modules/names.go b/edge/pkg/common/modules/names.go
index fa77bb205..47c085a0c 100644
--- a/edge/pkg/common/modules/names.go
+++ b/edge/pkg/common/modules/names.go
@@ -13,4 +13,6 @@ const (
DeviceTwinModuleName = "twin"
// EdgeHubModuleName name
EdgeHubModuleName = "websocket"
+ // MetaManagerModuleName metamanager module name
+ MetaManagerModuleName = "metamanager"
)
diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go
index 8788e2b97..5af9d4a70 100644
--- a/edge/pkg/edged/edged.go
+++ b/edge/pkg/edged/edged.go
@@ -115,7 +115,6 @@ import (
edgedutil "github.com/kubeedge/kubeedge/edge/pkg/edged/util"
"github.com/kubeedge/kubeedge/edge/pkg/edged/util/record"
csiplugin "github.com/kubeedge/kubeedge/edge/pkg/edged/volume/csi"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1"
"github.com/kubeedge/kubeedge/pkg/version"
@@ -1149,7 +1148,7 @@ func (e *edged) syncPod() {
//when starting, send msg to metamanager once to get existing pods
info := model.NewMessage("").BuildRouter(e.Name(), e.Group(), e.namespace+"/"+model.ResourceTypePod,
model.QueryOperation)
- beehiveContext.Send(metamanager.MetaManagerModuleName, *info)
+ beehiveContext.Send(modules.MetaManagerModuleName, *info)
for {
select {
case <-beehiveContext.Done():
@@ -1178,7 +1177,7 @@ func (e *edged) syncPod() {
klog.V(4).Infof("result content is %s", result.Content)
switch resType {
case model.ResourceTypePod:
- if op == model.ResponseOperation && resID == "" && result.GetSource() == metamanager.MetaManagerModuleName {
+ if op == model.ResponseOperation && resID == "" && result.GetSource() == modules.MetaManagerModuleName {
err := e.handlePodListFromMetaManager(content)
if err != nil {
klog.Errorf("handle podList failed: %v", err)
diff --git a/edge/pkg/edged/status/status_manager_test.go b/edge/pkg/edged/status/status_manager_test.go
index 09a4f780f..a40bb0e3c 100644
--- a/edge/pkg/edged/status/status_manager_test.go
+++ b/edge/pkg/edged/status/status_manager_test.go
@@ -38,11 +38,11 @@ func init() {
beehiveContext.InitContext([]string{common.MsgCtxTypeChannel})
metaManager := &common.ModuleInfo{
- ModuleName: "metaManager",
+ ModuleName: modules.MetaManagerModuleName,
ModuleType: common.MsgCtxTypeChannel,
}
beehiveContext.AddModule(metaManager)
- beehiveContext.AddModuleGroup("metaManager", "meta")
+ beehiveContext.AddModuleGroup(modules.MetaManagerModuleName, "meta")
edgeHub := &common.ModuleInfo{
ModuleName: "websocket",
@@ -95,7 +95,7 @@ func TestUpdatePodStatusSucceed(t *testing.T) {
newReason := "test reason"
manager.SetPodStatus(testPod, getPodStatus(newReason))
- msg, err := beehiveContext.Receive("metaManager")
+ msg, err := beehiveContext.Receive(modules.MetaManagerModuleName)
if err != nil {
t.Fatalf("receive message error: %v", err)
}
@@ -126,7 +126,7 @@ func TestUpdatePodStatusFailure(t *testing.T) {
oldReason := "test reason"
manager.SetPodStatus(testPod, getPodStatus(oldReason))
- msg, err := beehiveContext.Receive("metaManager")
+ msg, err := beehiveContext.Receive(modules.MetaManagerModuleName)
if err != nil {
t.Fatalf("receive message error: %v", err)
}
diff --git a/edge/pkg/edgehub/edgehub.go b/edge/pkg/edgehub/edgehub.go
index 7b4452458..f302e19f5 100644
--- a/edge/pkg/edgehub/edgehub.go
+++ b/edge/pkg/edgehub/edgehub.go
@@ -16,8 +16,6 @@ import (
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1"
)
-var HasTLSTunnelCerts = make(chan bool, 1)
-
//EdgeHub defines edgehub object structure
type EdgeHub struct {
certManager certificate.CertManager
@@ -30,7 +28,21 @@ type EdgeHub struct {
var _ core.Module = (*EdgeHub)(nil)
+var certSync map[string]chan bool
+
+func GetCertSyncChannel() map[string]chan bool {
+ return certSync
+}
+
+func NewCertSyncChannel() map[string]chan bool {
+ certSync = make(map[string]chan bool, 2)
+ certSync[modules.EdgeStreamModuleName] = make(chan bool, 1)
+ certSync[modules.MetaManagerModuleName] = make(chan bool, 1)
+ return certSync
+}
+
func newEdgeHub(enable bool) *EdgeHub {
+ NewCertSyncChannel()
return &EdgeHub{
enable: enable,
reconnectChan: make(chan struct{}),
@@ -65,9 +77,10 @@ func (eh *EdgeHub) Enable() bool {
func (eh *EdgeHub) Start() {
eh.certManager = certificate.NewCertManager(config.Config.EdgeHub, config.Config.NodeName)
eh.certManager.Start()
-
- HasTLSTunnelCerts <- true
- close(HasTLSTunnelCerts)
+ for _, v := range GetCertSyncChannel() {
+ v <- true
+ close(v)
+ }
go eh.ifRotationDone()
diff --git a/edge/pkg/edgestream/edgestream.go b/edge/pkg/edgestream/edgestream.go
index b6800c738..aac9edfc3 100644
--- a/edge/pkg/edgestream/edgestream.go
+++ b/edge/pkg/edgestream/edgestream.go
@@ -75,7 +75,7 @@ func (e *edgestream) Start() {
Path: "/v1/kubeedge/connect",
}
// TODO: Will improve in the future
- if ok := <-edgehub.HasTLSTunnelCerts; !ok {
+ if ok := <-edgehub.GetCertSyncChannel()[e.Name()]; !ok {
klog.Exitf("Failed to find cert key pair")
}
diff --git a/edge/pkg/metamanager/client/configmap.go b/edge/pkg/metamanager/client/configmap.go
index ea47860c1..ddc329f5c 100644
--- a/edge/pkg/metamanager/client/configmap.go
+++ b/edge/pkg/metamanager/client/configmap.go
@@ -9,7 +9,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
// ConfigMapsGetter has a method to return a ConfigMapInterface.
@@ -63,7 +62,7 @@ func (c *configMaps) Get(name string) (*api.ConfigMap, error) {
return nil, fmt.Errorf("parse message to configmap failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handleConfigMapFromMetaDB(content)
}
return handleConfigMapFromMetaManager(content)
diff --git a/edge/pkg/metamanager/client/metaclient.go b/edge/pkg/metamanager/client/metaclient.go
index 3a71e6828..8ef17d1e3 100644
--- a/edge/pkg/metamanager/client/metaclient.go
+++ b/edge/pkg/metamanager/client/metaclient.go
@@ -8,7 +8,7 @@ import (
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
"github.com/kubeedge/beehive/pkg/core/model"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
+ "github.com/kubeedge/kubeedge/edge/pkg/common/modules"
)
var (
@@ -102,7 +102,7 @@ func (s *send) SendSync(message *model.Message) (*model.Message, error) {
var resp model.Message
retries := 0
err = wait.Poll(syncPeriod, syncMsgRespTimeout, func() (bool, error) {
- resp, err = beehiveContext.SendSync(metamanager.MetaManagerModuleName, *message, syncMsgRespTimeout)
+ resp, err = beehiveContext.SendSync(modules.MetaManagerModuleName, *message, syncMsgRespTimeout)
retries++
if err == nil {
klog.V(4).Infof("send sync message %s succeed and response: %v", message.GetResource(), resp)
@@ -118,7 +118,7 @@ func (s *send) SendSync(message *model.Message) (*model.Message, error) {
}
func (s *send) Send(message *model.Message) {
- beehiveContext.Send(metamanager.MetaManagerModuleName, *message)
+ beehiveContext.Send(modules.MetaManagerModuleName, *message)
}
func SetSyncPeriod(time time.Duration) {
diff --git a/edge/pkg/metamanager/client/node.go b/edge/pkg/metamanager/client/node.go
index dc4d19870..9225c5256 100644
--- a/edge/pkg/metamanager/client/node.go
+++ b/edge/pkg/metamanager/client/node.go
@@ -9,7 +9,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
//NodesGetter to get node interface
@@ -68,7 +67,7 @@ func (c *nodes) Get(name string) (*api.Node, error) {
return nil, fmt.Errorf("parse message to node failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handleNodeFromMetaDB(content)
}
return handleNodeFromMetaManager(content)
diff --git a/edge/pkg/metamanager/client/persistentvolume.go b/edge/pkg/metamanager/client/persistentvolume.go
index 954f18c68..594eb0514 100644
--- a/edge/pkg/metamanager/client/persistentvolume.go
+++ b/edge/pkg/metamanager/client/persistentvolume.go
@@ -10,7 +10,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
// PersistentVolumesGetter is interface to get client PersistentVolumes
@@ -63,7 +62,7 @@ func (c *persistentvolumes) Get(name string, options metav1.GetOptions) (*api.Pe
return nil, fmt.Errorf("parse message to persistentvolume failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handlePersistentVolumeFromMetaDB(content)
}
return handlePersistentVolumeFromMetaManager(content)
diff --git a/edge/pkg/metamanager/client/persistentvolumeclaim.go b/edge/pkg/metamanager/client/persistentvolumeclaim.go
index d9fb3ce6a..c7002de14 100644
--- a/edge/pkg/metamanager/client/persistentvolumeclaim.go
+++ b/edge/pkg/metamanager/client/persistentvolumeclaim.go
@@ -10,7 +10,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
// PersistentVolumeClaimsGetter is interface to get client PersistentVolumeClaims
@@ -63,7 +62,7 @@ func (c *persistentvolumeclaims) Get(name string, options metav1.GetOptions) (*a
return nil, fmt.Errorf("parse message to persistentvolumeclaim failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handlePersistentVolumeClaimFromMetaDB(content)
}
return handlePersistentVolumeClaimFromMetaManager(content)
diff --git a/edge/pkg/metamanager/client/secret.go b/edge/pkg/metamanager/client/secret.go
index 2fe9388ee..17ed48d42 100644
--- a/edge/pkg/metamanager/client/secret.go
+++ b/edge/pkg/metamanager/client/secret.go
@@ -9,7 +9,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
//SecretsGetter is interface to get client secrets
@@ -62,7 +61,7 @@ func (c *secrets) Get(name string) (*api.Secret, error) {
return nil, fmt.Errorf("parse message to secret failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handleSecretFromMetaDB(content)
}
//else
diff --git a/edge/pkg/metamanager/client/serviceaccount.go b/edge/pkg/metamanager/client/serviceaccount.go
index 0a840db57..fe1bfc5d6 100644
--- a/edge/pkg/metamanager/client/serviceaccount.go
+++ b/edge/pkg/metamanager/client/serviceaccount.go
@@ -10,7 +10,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
// ServiceAccountTokenGetter is interface to get client service account token
@@ -48,7 +47,7 @@ func (c *serviceAccountToken) GetServiceAccountToken(namespace string, name stri
return nil, fmt.Errorf("marshal message to serviceaccount token failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handleServiceAccountTokenFromMetaDB(content)
}
return handleServiceAccountTokenFromMetaManager(content)
diff --git a/edge/pkg/metamanager/client/volumeattachment.go b/edge/pkg/metamanager/client/volumeattachment.go
index a080da40b..952b99a17 100644
--- a/edge/pkg/metamanager/client/volumeattachment.go
+++ b/edge/pkg/metamanager/client/volumeattachment.go
@@ -10,7 +10,6 @@ import (
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager"
)
// VolumeAttachmentsGetter is interface to get client VolumeAttachments
@@ -75,7 +74,7 @@ func (c *volumeattachments) Get(name string, options metav1.GetOptions) (*api.Vo
return nil, fmt.Errorf("parse message to volumeattachment failed, err: %v", err)
}
- if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == metamanager.MetaManagerModuleName {
+ if msg.GetOperation() == model.ResponseOperation && msg.GetSource() == modules.MetaManagerModuleName {
return handleVolumeAttachmentFromMetaDB(content)
}
return handleVolumeAttachmentFromMetaManager(content)
diff --git a/edge/pkg/metamanager/metamanager.go b/edge/pkg/metamanager/metamanager.go
index ec776edca..1b0a47930 100644
--- a/edge/pkg/metamanager/metamanager.go
+++ b/edge/pkg/metamanager/metamanager.go
@@ -16,11 +16,6 @@ import (
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1"
)
-//constant metamanager module name
-const (
- MetaManagerModuleName = "metaManager"
-)
-
type metaManager struct {
enable bool
}
@@ -53,7 +48,7 @@ func initDBTable(module core.Module) {
}
func (*metaManager) Name() string {
- return MetaManagerModuleName
+ return modules.MetaManagerModuleName
}
func (*metaManager) Group() string {
diff --git a/edge/pkg/metamanager/metamanager_test.go b/edge/pkg/metamanager/metamanager_test.go
index 3dd36d2e2..ae4636e57 100644
--- a/edge/pkg/metamanager/metamanager_test.go
+++ b/edge/pkg/metamanager/metamanager_test.go
@@ -31,7 +31,7 @@ var metaModule core.Module
func init() {
beehiveContext.InitContext([]string{common.MsgCtxTypeChannel})
add := &common.ModuleInfo{
- ModuleName: MetaManagerModuleName,
+ ModuleName: commodule.MetaManagerModuleName,
ModuleType: common.MsgCtxTypeChannel,
}
beehiveContext.AddModule(add)
@@ -41,7 +41,7 @@ func TestNameAndGroup(t *testing.T) {
modules := core.GetModules()
core.Register(&metaManager{enable: true})
for name, module := range modules {
- if name == MetaManagerModuleName {
+ if name == commodule.MetaManagerModuleName {
metaModule = module.GetModule()
break
}
@@ -51,8 +51,8 @@ func TestNameAndGroup(t *testing.T) {
t.Errorf("failed to register to beehive")
return
}
- if MetaManagerModuleName != metaModule.Name() {
- t.Errorf("Name of module is not correct wanted: %v and got: %v", MetaManagerModuleName, metaModule.Name())
+ if commodule.MetaManagerModuleName != metaModule.Name() {
+ t.Errorf("Name of module is not correct wanted: %v and got: %v", commodule.MetaManagerModuleName, metaModule.Name())
return
}
if commodule.MetaGroup != metaModule.Group() {
diff --git a/edge/pkg/metamanager/metaserver/config/config.go b/edge/pkg/metamanager/metaserver/config/config.go
index 071ab55fd..976d29709 100644
--- a/edge/pkg/metamanager/metaserver/config/config.go
+++ b/edge/pkg/metamanager/metaserver/config/config.go
@@ -17,9 +17,10 @@ type Configure struct {
func InitConfigure(c *v1alpha1.MetaServer) {
once.Do(func() {
- Config.Enable = c.Enable
- Config.Server = c.Server
- // so edgehub must register before metamanager
- Config.NodeName = edgehubconfig.Config.NodeName
+ Config = Configure{
+ MetaServer: *c,
+ // so edgehub must register before metamanager
+ NodeName: edgehubconfig.Config.NodeName,
+ }
})
}
diff --git a/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook/hooks.go b/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook/hooks.go
index 7f9579b87..9d8bce185 100644
--- a/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook/hooks.go
+++ b/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook/hooks.go
@@ -47,6 +47,7 @@ func Trigger(e watch.Event) {
return
}
gvr, ns, name := metaserver.ParseKey(key)
+ //TODO why not lock hooks ?
for _, hook := range hooks {
compGVR, compNS, compName, compRev := true, true, true, true
if !hook.GetGVR().Empty() {
diff --git a/edge/pkg/metamanager/metaserver/kubernetes/storage/storage.go b/edge/pkg/metamanager/metaserver/kubernetes/storage/storage.go
index c8ff92a2f..4579092b5 100644
--- a/edge/pkg/metamanager/metaserver/kubernetes/storage/storage.go
+++ b/edge/pkg/metamanager/metaserver/kubernetes/storage/storage.go
@@ -20,6 +20,7 @@ import (
"k8s.io/klog/v2"
"github.com/kubeedge/kubeedge/cloud/pkg/dynamiccontroller/application"
+ metaserverconfig "github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/config"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator"
"github.com/kubeedge/kubeedge/pkg/metaserver"
@@ -82,11 +83,15 @@ func (r *REST) Get(ctx context.Context, name string, options *metav1.GetOptions)
path := info.Path
// try remote cloud
obj, err := func() (runtime.Object, error) {
- app := r.Agent.Generate(ctx, application.Get, *options, nil)
- err := r.Agent.Apply(app)
+ app, err := r.Agent.Generate(ctx, application.Get, *options, nil)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, err
+ }
+ err = r.Agent.Apply(app)
defer app.Close()
if err != nil {
- klog.Errorf("[metaserver/reststorage] failed to get obj from cloud, %v", err)
+ klog.Errorf("[metaserver/reststorage] failed to get obj from cloud: %v", err)
return nil, err
}
var obj = new(unstructured.Unstructured)
@@ -100,7 +105,7 @@ func (r *REST) Get(ctx context.Context, name string, options *metav1.GetOptions)
return obj, nil
}()
// try local
- if err != nil {
+ if err != nil && metaserverconfig.Config.AutonomyWithoutAuthorization {
obj, err = r.Store.Get(ctx, "", options) // name is needless, we get all key information from ctx
if err != nil {
return nil, errors.NewNotFound(schema.GroupResource{Group: info.APIGroup, Resource: info.Resource}, info.Name)
@@ -114,8 +119,12 @@ func (r *REST) List(ctx context.Context, options *metainternalversion.ListOption
path := info.Path
// try remote cloud
list, err := func() (runtime.Object, error) {
- app := r.Agent.Generate(ctx, application.List, *options, nil)
- err := r.Agent.Apply(app)
+ app, err := r.Agent.Generate(ctx, application.List, *options, nil)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, err
+ }
+ err = r.Agent.Apply(app)
defer app.Close()
if err != nil {
return nil, err
@@ -132,7 +141,9 @@ func (r *REST) List(ctx context.Context, options *metainternalversion.ListOption
// try local if error occurs
if err != nil {
- klog.Warningf("[metaserver/reststorage] failed to list obj from cloud, %v; try local", err)
+ if !metaserverconfig.Config.AutonomyWithoutAuthorization {
+ return nil, err
+ }
list, err = r.Store.List(ctx, options)
if err != nil {
return nil, err
@@ -148,11 +159,15 @@ func (r *REST) Watch(ctx context.Context, options *metainternalversion.ListOptio
path := info.Path
// try remote cloud
_, err := func() (runtime.Object, error) {
- app := r.Agent.Generate(ctx, application.Watch, *options, nil)
- err := r.Agent.Apply(app)
+ app, err := r.Agent.Generate(ctx, application.Watch, *options, nil)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, err
+ }
+ err = r.Agent.Apply(app)
defer app.Close()
if err != nil {
- klog.Errorf("[metaserver/reststorage] failed to apply for a watch listener from cloud, %v", err)
+ klog.Errorf("[metaserver/reststorage] failed to apply for a watch listener from cloud: %v", err)
return nil, errors.NewInternalError(err)
}
klog.Infof("[metaserver/reststorage] successfully apply for a watch listener (%v) through cloud", path)
@@ -168,11 +183,15 @@ func (r *REST) Watch(ctx context.Context, options *metainternalversion.ListOptio
func (r *REST) Create(ctx context.Context, obj runtime.Object, createValidation rest.ValidateObjectFunc, options *metav1.CreateOptions) (runtime.Object, error) {
obj, err := func() (runtime.Object, error) {
- app := r.Agent.Generate(ctx, application.Create, *options, obj)
- err := r.Agent.Apply(app)
+ app, err := r.Agent.Generate(ctx, application.Create, *options, obj)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, err
+ }
+ err = r.Agent.Apply(app)
defer app.Close()
if err != nil {
- klog.Errorf("[metaserver/reststorage] failed to create obj, %v", err)
+ klog.Errorf("[metaserver/reststorage] failed to create obj: %v", err)
return nil, err
}
@@ -192,8 +211,12 @@ func (r *REST) Create(ctx context.Context, obj runtime.Object, createValidation
func (r *REST) Delete(ctx context.Context, name string, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions) (runtime.Object, bool, error) {
key, _ := metaserver.KeyFuncReq(ctx, "")
- app := r.Agent.Generate(ctx, application.Delete, options, nil)
- err := r.Agent.Apply(app)
+ app, err := r.Agent.Generate(ctx, application.Delete, options, nil)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, false, err
+ }
+ err = r.Agent.Apply(app)
defer app.Close()
if err != nil {
klog.Errorf("[metaserver/reststorage] failed to delete (%v) through cloud", key)
@@ -212,9 +235,13 @@ func (r *REST) Update(ctx context.Context, name string, objInfo rest.UpdatedObje
reqInfo, _ := apirequest.RequestInfoFrom(ctx)
var app *application.Application
if reqInfo.Subresource == "status" {
- app = r.Agent.Generate(ctx, application.UpdateStatus, options, obj)
+ app, err = r.Agent.Generate(ctx, application.UpdateStatus, options, obj)
} else {
- app = r.Agent.Generate(ctx, application.Update, options, obj)
+ app, err = r.Agent.Generate(ctx, application.Update, options, obj)
+ }
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, false, err
}
defer app.Close()
if err := r.Agent.Apply(app); err != nil {
@@ -228,7 +255,11 @@ func (r *REST) Update(ctx context.Context, name string, objInfo rest.UpdatedObje
}
func (r *REST) Patch(ctx context.Context, pi application.PatchInfo) (runtime.Object, error) {
- app := r.Agent.Generate(ctx, application.Patch, pi, nil)
+ app, err := r.Agent.Generate(ctx, application.Patch, pi, nil)
+ if err != nil {
+ klog.Errorf("[metaserver/reststorage] failed to generate application: %v", err)
+ return nil, err
+ }
defer app.Close()
if err := r.Agent.Apply(app); err != nil {
return nil, err
diff --git a/edge/pkg/metamanager/metaserver/server.go b/edge/pkg/metamanager/metaserver/server.go
index 342bf183b..0327ba17b 100644
--- a/edge/pkg/metamanager/metaserver/server.go
+++ b/edge/pkg/metamanager/metaserver/server.go
@@ -2,8 +2,12 @@ package metaserver
import (
"context"
+ "crypto/tls"
+ "crypto/x509"
+ "encoding/pem"
"fmt"
"net/http"
+ "os"
"time"
"k8s.io/apimachinery/pkg/api/errors"
@@ -17,9 +21,13 @@ import (
apirequest "k8s.io/apiserver/pkg/endpoints/request"
"k8s.io/apiserver/pkg/server"
genericfilters "k8s.io/apiserver/pkg/server/filters"
+ certutil "k8s.io/client-go/util/cert"
"k8s.io/klog/v2"
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
+ "github.com/kubeedge/kubeedge/edge/pkg/common/modules"
+ "github.com/kubeedge/kubeedge/edge/pkg/edgehub"
metaserverconfig "github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/config"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/handlerfactory"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/kubernetes/serializer"
@@ -45,12 +53,72 @@ func NewMetaServer() *MetaServer {
return &ls
}
+func createTLSConfig() tls.Config {
+ ca, err := os.ReadFile(metaserverconfig.Config.TLSCaFile)
+ if err == nil {
+ block, _ := pem.Decode(ca)
+ ca = block.Bytes
+ klog.Info("Succeed in loading CA certificate from local directory")
+ }
+ pool := x509.NewCertPool()
+ ok := pool.AppendCertsFromPEM(pem.EncodeToMemory(&pem.Block{Type: certutil.CertificateBlockType, Bytes: ca}))
+ if !ok {
+ panic(fmt.Errorf("fail to load ca content"))
+ }
+ cert, err := os.ReadFile(metaserverconfig.Config.TLSCertFile)
+ if err == nil {
+ block, _ := pem.Decode(cert)
+ cert = block.Bytes
+ klog.Info("Succeed in loading certificate from local directory")
+ }
+ key, err := os.ReadFile(metaserverconfig.Config.TLSPrivateKeyFile)
+ if err == nil {
+ block, _ := pem.Decode(key)
+ key = block.Bytes
+ klog.Info("Succeed in loading private key from local directory")
+ }
+
+ certificate, err := tls.X509KeyPair(pem.EncodeToMemory(&pem.Block{Type: certutil.CertificateBlockType, Bytes: cert}), pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: key}))
+ if err != nil {
+ panic(err)
+ }
+ return tls.Config{
+ ClientCAs: pool,
+ ClientAuth: tls.VerifyClientCertIfGiven,
+ Certificates: []tls.Certificate{certificate},
+ MinVersion: tls.VersionTLS12,
+ }
+}
+
+// getCurrent returns current meta server certificate
+func (ls *MetaServer) getCurrent() (*tls.Certificate, error) {
+ cert, err := tls.LoadX509KeyPair(metaserverconfig.Config.TLSCertFile, metaserverconfig.Config.TLSPrivateKeyFile)
+ if err != nil {
+ return nil, err
+ }
+ certs, err := x509.ParseCertificates(cert.Certificate[0])
+ if err != nil {
+ return nil, fmt.Errorf("unable to parse certificate data: %v", err)
+ }
+ cert.Leaf = certs[0]
+ return &cert, nil
+}
+
func (ls *MetaServer) Start(stopChan <-chan struct{}) {
+ _, err := ls.getCurrent()
+ if err != nil {
+ // wait for cert created
+ klog.Infof("[metaserver]waiting for cert created")
+ <-edgehub.GetCertSyncChannel()[modules.MetaManagerModuleName]
+ }
+
h := ls.BuildBasicHandler()
h = BuildHandlerChain(h, ls)
+ tlsConfig := createTLSConfig()
s := http.Server{
- Addr: metaserverconfig.Config.Server,
- Handler: h,
+ Addr: metaserverconfig.Config.Server,
+ Handler: h,
+ TLSConfig: &tlsConfig,
}
go func() {
@@ -64,13 +132,13 @@ func (ls *MetaServer) Start(stopChan <-chan struct{}) {
}()
klog.Infof("[metaserver]start to listen and server at %v", s.Addr)
- utilruntime.HandleError(s.ListenAndServe())
+ utilruntime.HandleError(s.ListenAndServeTLS("", ""))
// When the MetaServer stops abnormally, other module services are stopped at the same time.
beehiveContext.Cancel()
}
func (ls *MetaServer) BuildBasicHandler() http.Handler {
- h := http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
+ return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
ctx := req.Context()
reqInfo, ok := apirequest.RequestInfoFrom(ctx)
//klog.Infof("[metaserver]get a req(%v)(%v)", reqInfo.Path, reqInfo.Verb)
@@ -92,14 +160,13 @@ func (ls *MetaServer) BuildBasicHandler() http.Handler {
default:
err := fmt.Errorf("unsupported req verb")
responsewriters.ErrorNegotiated(errors.NewInternalError(err), ls.NegotiatedSerializer, schema.GroupVersion{}, w, req)
- return
}
- } else {
- err := fmt.Errorf("not a resource req")
- responsewriters.ErrorNegotiated(errors.NewInternalError(err), ls.NegotiatedSerializer, schema.GroupVersion{}, w, req)
+ return
}
+
+ err := fmt.Errorf("not a resource req")
+ responsewriters.ErrorNegotiated(errors.NewInternalError(err), ls.NegotiatedSerializer, schema.GroupVersion{}, w, req)
})
- return h
}
func BuildHandlerChain(handler http.Handler, ls *MetaServer) http.Handler {
@@ -107,10 +174,21 @@ func BuildHandlerChain(handler http.Handler, ls *MetaServer) http.Handler {
LegacyAPIGroupPrefixes: sets.NewString(server.DefaultLegacyAPIPrefix),
}
- requestInfoResolver := &apirequest.RequestInfoFactory{}
-
handler = genericfilters.WithWaitGroup(handler, ls.LongRunningFunc, ls.HandlerChainWaitGroup)
handler = genericapifilters.WithRequestInfo(handler, server.NewRequestInfoResolver(cfg))
- handler = genericfilters.WithPanicRecovery(handler, requestInfoResolver)
+ handler = genericfilters.WithPanicRecovery(handler, &apirequest.RequestInfoFactory{})
+ handler = CheckAuthorizationHeader(handler)
return handler
}
+
+func CheckAuthorizationHeader(handler http.Handler) http.Handler {
+ return http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
+ token := request.Header.Get(commontypes.AuthorizationKey)
+ if len(token) == 0 {
+ http.Error(writer, "header Authorization is missing", http.StatusNetworkAuthenticationRequired)
+ return
+ }
+ request = request.WithContext(context.WithValue(request.Context(), commontypes.AuthorizationKey, token))
+ handler.ServeHTTP(writer, request)
+ })
+}
diff --git a/edge/pkg/metamanager/process.go b/edge/pkg/metamanager/process.go
index 19eabcc5e..79f840ee3 100644
--- a/edge/pkg/metamanager/process.go
+++ b/edge/pkg/metamanager/process.go
@@ -39,7 +39,7 @@ func feedbackError(err error, info string, request model.Message) {
if err != nil {
errInfo = fmt.Sprintf(info+": %v", err)
}
- errResponse := model.NewErrorMessage(&request, errInfo).SetRoute(MetaManagerModuleName, request.GetGroup())
+ errResponse := model.NewErrorMessage(&request, errInfo).SetRoute(modules.MetaManagerModuleName, request.GetGroup())
if request.GetSource() == modules.EdgedModuleName {
sendToEdged(errResponse, request.IsSync())
} else {
@@ -262,7 +262,7 @@ func (m *metaManager) processQuery(message model.Message) {
m.processRemoteQuery(message)
} else {
resp := message.NewRespByMessage(&message, *metas)
- resp.SetRoute(MetaManagerModuleName, resp.GetGroup())
+ resp.SetRoute(modules.MetaManagerModuleName, resp.GetGroup())
sendToEdged(resp, message.IsSync())
}
return
@@ -279,7 +279,7 @@ func (m *metaManager) processQuery(message model.Message) {
feedbackError(err, "Error to query meta in DB", message)
} else {
resp := message.NewRespByMessage(&message, *metas)
- resp.SetRoute(MetaManagerModuleName, resp.GetGroup())
+ resp.SetRoute(modules.MetaManagerModuleName, resp.GetGroup())
sendToEdged(resp, message.IsSync())
}
}
diff --git a/edge/pkg/metamanager/process_test.go b/edge/pkg/metamanager/process_test.go
index 2e4e888df..dca05e2d7 100644
--- a/edge/pkg/metamanager/process_test.go
+++ b/edge/pkg/metamanager/process_test.go
@@ -61,7 +61,7 @@ func init() {
beehiveContext.InitContext([]string{common.MsgCtxTypeChannel})
add := &common.ModuleInfo{
- ModuleName: MetaManagerModuleName,
+ ModuleName: modules.MetaManagerModuleName,
ModuleType: common.MsgCtxTypeChannel,
}
beehiveContext.AddModule(add)
@@ -96,7 +96,7 @@ func TestProcessInsert(t *testing.T) {
//SaveMeta Failed, feedbackError SendToCloud
ormerMock.EXPECT().Insert(gomock.Any()).Return(int64(1), errFailedDBOperation).Times(1)
- msg := model.NewMessage("").BuildRouter(MetaManagerModuleName, GroupResource, model.ResourceTypePodStatus, model.InsertOperation)
+ msg := model.NewMessage("").BuildRouter(modules.MetaManagerModuleName, GroupResource, model.ResourceTypePodStatus, model.InsertOperation)
meta.processInsert(*msg)
//beehiveContext.Send(MetaManagerModuleName, *msg)
message, err := beehiveContext.Receive(ModuleNameEdgeHub)
@@ -337,7 +337,7 @@ func TestProcessResponse(t *testing.T) {
//jsonMarshall fail
msg := model.NewMessage("").BuildRouter(ModuleNameEdged, GroupResource, model.ResourceTypePodStatus, model.ResponseOperation).FillBody(make(chan int))
- beehiveContext.Send(MetaManagerModuleName, *msg)
+ beehiveContext.Send(modules.MetaManagerModuleName, *msg)
meta.processResponse(*msg)
message, _ := beehiveContext.Receive(ModuleNameEdged)
t.Run("MarshallFail", func(t *testing.T) {
diff --git a/edge/test/integration/metaserver/metaserver_suite_test.go b/edge/test/integration/metaserver/metaserver_suite_test.go
index f1adc8a51..f006cc3f9 100644
--- a/edge/test/integration/metaserver/metaserver_suite_test.go
+++ b/edge/test/integration/metaserver/metaserver_suite_test.go
@@ -25,6 +25,10 @@ func TestEdgecoreMetaServer(t *testing.T) {
c.Modules.Edged.HostnameOverride = cfg.NodeID
c.Modules.MetaManager.Enable = true
c.Modules.MetaManager.MetaServer.Enable = true
+ c.Modules.MetaManager.MetaServer.AutonomyWithoutAuthorization = true
+ c.Modules.MetaManager.MetaServer.TLSCaFile = "/tmp/edgecore/rootCA.crt"
+ c.Modules.MetaManager.MetaServer.TLSCertFile = "/tmp/edgecore/kubeedge.crt"
+ c.Modules.MetaManager.MetaServer.TLSPrivateKeyFile = "/tmp/edgecore/kubeedge.key"
Expect(utils.CfgToFile(c)).Should(BeNil())
Expect(utils.StartEdgeCore()).Should(BeNil())
diff --git a/edge/test/integration/metaserver/metaserver_test.go b/edge/test/integration/metaserver/metaserver_test.go
index 09fb2eb07..120a9a7a9 100644
--- a/edge/test/integration/metaserver/metaserver_test.go
+++ b/edge/test/integration/metaserver/metaserver_test.go
@@ -1,6 +1,7 @@
package metaserver
import (
+ "crypto/tls"
"net/http"
. "github.com/onsi/ginkgo/v2"
@@ -42,11 +43,17 @@ var _ = Describe("Test MetaServer", func() {
//"Watch with bad method": {"POST", "/" + prefix + "/" + testGroupVersion.Group + "/" + testGroupVersion.Version + "/watch/namespaces/ns/simples/", http.StatusMethodNotAllowed},
//"Watch param with bad method": {"POST", "/" + prefix + "/" + testGroupVersion.Group + "/" + testGroupVersion.Version + "/namespaces/ns-foo/simples?watch=true", http.StatusMethodNotAllowed},
}
- client := http.Client{}
- url := "http://127.0.0.1:10550"
+ client := http.Client{
+ Transport: &http.Transport{
+ TLSClientConfig: &tls.Config{
+ InsecureSkipVerify: true},
+ },
+ }
+ url := "https://127.0.0.1:10550"
for _, v := range cases {
request, err := http.NewRequest(v.Method, url+v.Path, nil)
Expect(err).Should(BeNil())
+ request.Header.Set("Authorization", "xxxxx")
response, err := client.Do(request)
Expect(err).Should(BeNil())
isEqual := v.Status == response.StatusCode
diff --git a/edge/test/test.go b/edge/test/test.go
index fedf82251..82d4919cd 100644
--- a/edge/test/test.go
+++ b/edge/test/test.go
@@ -122,7 +122,7 @@ func (tm *testManager) podHandler(w http.ResponseWriter, req *http.Request) {
ns = p.Namespace
}
msgReq := message.BuildMsg("resource", string(p.UID), "edgecontroller", ns+"/pod/"+p.Name, operation, p)
- beehiveContext.Send("metaManager", *msgReq)
+ beehiveContext.Send(modules.MetaManagerModuleName, *msgReq)
klog.Infof("send message to metaManager is %+v\n", msgReq)
}
}
@@ -183,7 +183,7 @@ func (tm *testManager) secretHandler(w http.ResponseWriter, req *http.Request) {
}
msgReq := message.BuildMsg("edgehub", string(p.UID), "test", "fakeNamespace/secret/"+string(p.UID), operation, p)
- beehiveContext.Send("metaManager", *msgReq)
+ beehiveContext.Send(modules.MetaManagerModuleName, *msgReq)
klog.Infof("send message to metaManager is %+v\n", msgReq)
}
}
@@ -213,7 +213,7 @@ func (tm *testManager) configmapHandler(w http.ResponseWriter, req *http.Request
}
msgReq := message.BuildMsg("edgehub", string(p.UID), "test", "fakeNamespace/configmap/"+string(p.UID), operation, p)
- beehiveContext.Send("metaManager", *msgReq)
+ beehiveContext.Send(modules.MetaManagerModuleName, *msgReq)
klog.Infof("send message to metaManager is %+v\n", msgReq)
}
}
diff --git a/pkg/apis/componentconfig/edgecore/v1alpha1/default.go b/pkg/apis/componentconfig/edgecore/v1alpha1/default.go
index 67bade00a..65d965b68 100644
--- a/pkg/apis/componentconfig/edgecore/v1alpha1/default.go
+++ b/pkg/apis/componentconfig/edgecore/v1alpha1/default.go
@@ -139,8 +139,12 @@ func NewDefaultEdgeCoreConfig() *EdgeCoreConfig {
ContextSendModule: metaconfig.ModuleNameEdgeHub,
RemoteQueryTimeout: constants.DefaultRemoteQueryTimeout,
MetaServer: &MetaServer{
- Enable: false,
- Server: constants.DefaultMetaServerAddr,
+ Enable: false,
+ AutonomyWithoutAuthorization: false,
+ Server: constants.DefaultMetaServerAddr,
+ TLSCaFile: constants.DefaultCAFile,
+ TLSCertFile: constants.DefaultCertFile,
+ TLSPrivateKeyFile: constants.DefaultKeyFile,
},
},
ServiceBus: &ServiceBus{
diff --git a/pkg/apis/componentconfig/edgecore/v1alpha1/types.go b/pkg/apis/componentconfig/edgecore/v1alpha1/types.go
index 41867041f..2c458ea40 100644
--- a/pkg/apis/componentconfig/edgecore/v1alpha1/types.go
+++ b/pkg/apis/componentconfig/edgecore/v1alpha1/types.go
@@ -399,8 +399,15 @@ type MetaManager struct {
}
type MetaServer struct {
- Enable bool `json:"enable"`
- Server string `json:"server"`
+ Enable bool `json:"enable"`
+ // AutonomyWithoutAuthorization is a switch to determine whether app can List/Watch
+ // meta data from local host db without authorization when the edge node is off-line.
+ // The default value is false, means won't be allowed.
+ AutonomyWithoutAuthorization bool `json:"autonomyWithoutAuthorization"`
+ Server string `json:"server"`
+ TLSCaFile string `json:"tlsCaFile"`
+ TLSCertFile string `json:"tlsCertFile"`
+ TLSPrivateKeyFile string `json:"tlsPrivateKeyFile"`
}
// ServiceBus indicates the ServiceBus module config
diff --git a/pkg/apis/componentconfig/meta/v1alpha1/types.go b/pkg/apis/componentconfig/meta/v1alpha1/types.go
index bab427c42..163cfe011 100644
--- a/pkg/apis/componentconfig/meta/v1alpha1/types.go
+++ b/pkg/apis/componentconfig/meta/v1alpha1/types.go
@@ -25,7 +25,7 @@ const (
ModuleNameServiceBus ModuleName = "servicebus"
// TODO @kadisi change websocket to edgehub
ModuleNameEdgeHub ModuleName = "websocket"
- ModuleNameMetaManager ModuleName = "metaManager"
+ ModuleNameMetaManager ModuleName = "metamanager"
ModuleNameEdged ModuleName = "edged"
ModuleNameTwin ModuleName = "twin"
ModuleNameDBTest ModuleName = "dbTest"