diff options
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" |
