kubeedge/cloud/pkg/dynamiccontroller/application/application.go
WillardHu 9c9f0c9d83 Fixed the golangci config and some code validation issues
Signed-off-by: WillardHu <wei.hu@daocloud.io>
2025-01-23 11:03:03 +08:00

397 lines
13 KiB
Go

package application
import (
"context"
"fmt"
"strings"
authorizationv1 "k8s.io/api/authorization/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/dynamic/dynamicinformer"
"k8s.io/client-go/kubernetes"
"k8s.io/klog/v2"
"k8s.io/kubernetes/cmd/kubeadm/app/constants"
"github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/cloud/pkg/common/client"
utilcontext "github.com/kubeedge/kubeedge/cloud/pkg/common/context"
"github.com/kubeedge/kubeedge/cloud/pkg/common/messagelayer"
"github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
"github.com/kubeedge/kubeedge/cloud/pkg/dynamiccontroller/config"
"github.com/kubeedge/kubeedge/cloud/pkg/dynamiccontroller/filter"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/pkg/metaserver"
passthrough "github.com/kubeedge/kubeedge/pkg/util/pass-through"
)
type Center struct {
HandlerCenter
messageLayer messagelayer.MessageLayer
dynamicClient dynamic.Interface
kubeClient kubernetes.Interface
}
func NewApplicationCenter(dynamicSharedInformerFactory dynamicinformer.DynamicSharedInformerFactory) *Center {
a := &Center{
HandlerCenter: NewHandlerCenter(dynamicSharedInformerFactory),
dynamicClient: client.GetDynamicClient(),
kubeClient: client.GetKubeClient(),
messageLayer: messagelayer.DynamicControllerMessageLayer(),
}
return a
}
// Process translate msg to application , process and send resp to edge
// TODO: upgrade to parallel process
func (c *Center) Process(msg model.Message) {
if strings.HasSuffix(msg.GetResource(), metaserver.WatchAppSync) {
if err := c.ProcessWatchSync(msg); err != nil {
klog.Errorf("failed to ProcessWatchSync: %v", err)
}
return
}
app, err := metaserver.MsgToApplication(msg)
if err != nil {
klog.Errorf("failed to translate msg to Application: %v", err)
return
}
klog.Infof("[metaserver/ApplicationCenter] get a Application %v", app.String())
if config.Config.EnableAuthorization {
nodeID, err := messagelayer.GetNodeID(msg)
if err != nil || nodeID != app.Nodename {
klog.Errorf("[metaserver/authorization]failed to process Application(%+v), %v", app, err)
return
}
}
if passthrough.IsPassThroughPath(app.Key, string(app.Verb)) {
resp, err := c.passThroughRequest(app)
if err != nil {
c.Response(app, msg.GetID(), metaserver.Rejected, err, nil)
klog.Errorf("[metaserver/passThrough]failed to process Application(%+v), %v", app, err)
return
}
c.Response(app, msg.GetID(), metaserver.Approved, nil, resp)
return
}
resp, err := c.ProcessApplication(app)
if err != nil {
c.Response(app, msg.GetID(), metaserver.Rejected, err, nil)
klog.Errorf("[metaserver/applicationCenter]failed to process Application(%+v), %v", app, err)
return
}
c.Response(app, msg.GetID(), metaserver.Approved, nil, resp)
klog.Infof("[metaserver/applicationCenter]successfully to process Application(%+v)", app)
}
// 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
// than request to apiserver directly. Listener will then continuously listen kube-api change events and
// push them to edge node.
func (c *Center) ProcessApplication(app *metaserver.Application) (interface{}, error) {
app.Status = metaserver.InProcessing
gvr, ns, name := metaserver.ParseKey(app.Key)
switch app.Verb {
case metaserver.List:
var option = new(metav1.ListOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
list, err := c.dynamicClient.Resource(gvr).Namespace(ns).List(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), *option)
if err != nil {
return nil, fmt.Errorf("get current list error: %v", err)
}
return list, nil
case metaserver.Watch:
if err := c.checkNodePermission(app); err != nil {
return nil, err
}
listener, err := applicationToListener(app)
if err != nil {
return nil, err
}
if err := c.HandlerCenter.AddListener(listener); err != nil {
return nil, fmt.Errorf("failed to add listener, %v", err)
}
return nil, nil
case metaserver.Get:
var option = new(metav1.GetOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
retObj, err := c.dynamicClient.Resource(gvr).Namespace(ns).Get(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), name, *option)
if err != nil {
return nil, err
}
return retObj, nil
case metaserver.Create:
var option = new(metav1.CreateOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
var obj = new(unstructured.Unstructured)
if err := app.ReqBodyTo(obj); err != nil {
return nil, err
}
var retObj interface{}
var err error
if app.Subresource == "" {
retObj, err = c.dynamicClient.Resource(gvr).Namespace(ns).Create(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), obj, *option)
} else {
retObj, err = c.dynamicClient.Resource(gvr).Namespace(ns).Create(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), obj, *option, app.Subresource)
}
if err != nil {
return nil, err
}
return retObj, err
case metaserver.Delete:
var option = new(metav1.DeleteOptions)
if err := app.OptionTo(&option); err != nil {
return nil, err
}
if err := c.dynamicClient.Resource(gvr).Namespace(ns).Delete(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), name, *option); err != nil {
return nil, err
}
return nil, nil
case metaserver.Update:
var option = new(metav1.UpdateOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
var obj = new(unstructured.Unstructured)
if err := app.ReqBodyTo(obj); err != nil {
return nil, err
}
var retObj interface{}
var err error
if app.Subresource == "" {
retObj, err = c.dynamicClient.Resource(gvr).Namespace(ns).Update(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), obj, *option)
} else {
retObj, err = c.dynamicClient.Resource(gvr).Namespace(ns).Update(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), obj, *option, app.Subresource)
}
if err != nil {
return nil, err
}
return retObj, nil
case metaserver.UpdateStatus:
var option = new(metav1.UpdateOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
var obj = new(unstructured.Unstructured)
if err := app.ReqBodyTo(obj); err != nil {
return nil, err
}
retObj, err := c.dynamicClient.Resource(gvr).Namespace(ns).UpdateStatus(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), obj, *option)
if err != nil {
return nil, err
}
return retObj, nil
case metaserver.Patch:
var pi = new(metaserver.PatchInfo)
if err := app.OptionTo(pi); err != nil {
return nil, err
}
retObj, err := c.dynamicClient.Resource(gvr).Namespace(ns).Patch(utilcontext.WithEdgeNode(context.TODO(), app.Nodename), pi.Name, pi.PatchType, pi.Data, pi.Options, pi.Subresources...)
if err != nil {
return nil, err
}
return retObj, nil
default:
return nil, fmt.Errorf("unsupported Application Verb type :%v", app.Verb)
}
}
func (c *Center) passThroughRequest(app *metaserver.Application) (interface{}, error) {
kubeClient, ok := c.kubeClient.(*kubernetes.Clientset)
if !ok {
return nil, fmt.Errorf("converting kubeClient to *kubernetes.Clientset type failed")
}
verb := strings.ToUpper(string(app.Verb))
return kubeClient.RESTClient().Verb(verb).AbsPath(app.Key).Body(app.ReqBody).Do(utilcontext.WithEdgeNode(context.TODO(), app.Nodename)).Raw()
}
// Response update application, generate and send resp message to edge
func (c *Center) Response(app *metaserver.Application, parentID string, status metaserver.ApplicationStatus, err error, respContent interface{}) {
app.Status = status
if err != nil {
apierr, ok := err.(apierrors.APIStatus)
if ok {
app.Error = apierrors.StatusError{ErrStatus: apierr.Status()}
} else {
app.Reason = err.Error()
}
}
if respContent != nil {
if app.Verb == metaserver.List || app.Verb == metaserver.Get {
filter.MessageFilter(respContent, app.Nodename)
}
app.RespBody = metaserver.ToBytes(respContent)
}
resource, err := messagelayer.BuildResource(app.Nodename, metaserver.Ignore, metaserver.ApplicationResource, metaserver.Ignore)
if err != nil {
klog.Warningf("built message resource failed with error: %s", err)
return
}
msg := model.NewMessage(parentID).
BuildRouter(modules.DynamicControllerModuleName, message.ResourceGroupName, resource, metaserver.ApplicationResp).
FillBody(app)
if err := c.messageLayer.Response(*msg); err != nil {
klog.Warningf("send message failed with error: %s, operation: %s, resource: %s", err, msg.GetOperation(), msg.GetResource())
return
}
klog.V(4).Infof("send message successfully, operation: %s, resource: %s", msg.GetOperation(), msg.GetResource())
}
// ProcessWatchSync process watch sync message
func (c *Center) ProcessWatchSync(msg model.Message) error {
nodeID, err := messagelayer.GetNodeID(msg)
if err != nil {
return err
}
applications, err := metaserver.MsgToApplications(msg)
if err != nil {
return fmt.Errorf("failed translate msg to Applications: %v", err)
}
addedWatchApp, removedListeners := c.getWatchDiff(applications, nodeID)
// gc already removed listeners
for _, listener := range removedListeners {
c.HandlerCenter.DeleteListener(listener)
}
failedWatchApp := make(map[string]metaserver.Application)
// add listener for new added watch app
for _, watchApp := range addedWatchApp {
if config.Config.EnableAuthorization && nodeID != watchApp.Nodename {
return fmt.Errorf("node name %q is not allowed", watchApp.Nodename)
}
err := c.processWatchApp(&watchApp)
if err != nil {
watchApp.Status = metaserver.Rejected
apiErr, ok := err.(apierrors.APIStatus)
if ok {
watchApp.Error = apierrors.StatusError{ErrStatus: apiErr.Status()}
} else {
watchApp.Reason = err.Error()
}
failedWatchApp[watchApp.ID] = watchApp
klog.Errorf("processWatchApp %s err: %v", watchApp.String(), err)
}
}
respMsg := model.NewMessage(msg.GetID()).
BuildRouter(modules.DynamicControllerModuleName, message.ResourceGroupName, msg.GetResource(), metaserver.ApplicationResp).
FillBody(failedWatchApp)
if err := c.messageLayer.Response(*respMsg); err != nil {
klog.Warningf("send message failed error: %s, operation: %s, resource: %s", err, respMsg.GetOperation(), respMsg.GetResource())
return err
}
return nil
}
func (c *Center) getWatchDiff(allWatchAppInEdge map[string]metaserver.Application,
nodeID string) ([]metaserver.Application, []*SelectorListener) {
listenerInCloud := c.HandlerCenter.GetListenersForNode(nodeID)
addedWatchApp := make([]metaserver.Application, 0)
for ID, app := range allWatchAppInEdge {
if _, exist := listenerInCloud[ID]; !exist {
addedWatchApp = append(addedWatchApp, app)
klog.Infof("added watch app %s", app.String())
}
}
removedListeners := make([]*SelectorListener, 0)
for ID, listener := range listenerInCloud {
if _, exist := allWatchAppInEdge[ID]; !exist {
removedListeners = append(removedListeners, listener)
klog.Infof("need removed listener %s", listener.id)
}
}
return addedWatchApp, removedListeners
}
func (c *Center) processWatchApp(watchApp *metaserver.Application) error {
if err := c.checkNodePermission(watchApp); err != nil {
return err
}
watchApp.Status = metaserver.InProcessing
listener, err := applicationToListener(watchApp)
if err != nil {
return err
}
if err := c.HandlerCenter.AddListener(listener); err != nil {
return fmt.Errorf("failed to add listener, %v", err)
}
return nil
}
func (c *Center) checkNodePermission(app *metaserver.Application) error {
if !config.Config.EnableAuthorization {
return nil
}
gvr, ns, name := metaserver.ParseKey(app.Key)
subjectAccessReview := &authorizationv1.SubjectAccessReview{
Spec: authorizationv1.SubjectAccessReviewSpec{
ResourceAttributes: &authorizationv1.ResourceAttributes{
Namespace: ns,
Verb: string(app.Verb),
Group: gvr.Group,
Version: gvr.Version,
Resource: gvr.Resource,
Subresource: app.Subresource,
Name: name,
},
User: constants.NodesUserPrefix + app.Nodename,
Groups: []string{constants.NodesGroup},
},
}
ret, err := c.kubeClient.AuthorizationV1().SubjectAccessReviews().Create(context.TODO(), subjectAccessReview, metav1.CreateOptions{})
if err != nil {
return fmt.Errorf("node %s permission check failed: %v", app.Nodename, err)
}
if !ret.Status.Allowed {
return fmt.Errorf("node %q is not allowed to access this resource", app.Nodename)
}
return nil
}
func applicationToListener(app *metaserver.Application) (*SelectorListener, error) {
var option = new(metav1.ListOptions)
if err := app.OptionTo(option); err != nil {
return nil, err
}
gvr, namespace, _ := metaserver.ParseKey(app.Key)
selector := NewSelector(option.LabelSelector, option.FieldSelector)
if namespace != "" {
selector.Field = fields.AndSelectors(selector.Field, fields.OneTermEqualSelector("metadata.namespace", namespace))
}
return NewSelectorListener(app.ID, app.Nodename, gvr, selector), nil
}