diff options
| author | edisonxiang <xiang.edison@gmail.com> | 2019-08-22 00:53:38 +0800 |
|---|---|---|
| committer | edisonxiang <xiang.edison@gmail.com> | 2019-09-04 10:09:15 +0800 |
| commit | 8939add58f4d560f55ed2b126483a04e355c32ff (patch) | |
| tree | 8ea9c970c508675460dbc64f130556cf6a5b1619 /edge | |
| parent | update required vendors (diff) | |
| download | kubeedge-8939add58f4d560f55ed2b126483a04e355c32ff.tar.gz | |
invoke func from edged
Diffstat (limited to 'edge')
| -rw-r--r-- | edge/pkg/edged/edged.go | 99 |
1 files changed, 99 insertions, 0 deletions
diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go index 8d7eae101..6e7ac2f2a 100644 --- a/edge/pkg/edged/edged.go +++ b/edge/pkg/edged/edged.go @@ -24,6 +24,7 @@ and made some variant package edged import ( + "bytes" "encoding/json" "fmt" "net" @@ -33,6 +34,8 @@ import ( "sync" "time" + "github.com/container-storage-interface/spec/lib/go/csi" + "github.com/golang/protobuf/jsonpb" cadvisorapi "github.com/google/cadvisor/info/v1" v1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -979,6 +982,15 @@ func (e *edged) syncPod() { klog.Infof("skip to handle secret with type response") continue } + case constants.CSIResourceTypeVolume: + klog.Infof("volume operation type: %s", op) + res, err := e.handleVolume(op, content) + if err != nil { + klog.Errorf("handle volume failed: %v", err) + } else { + resp := request.NewRespByMessage(&request, res) + e.context.SendResp(*resp) + } default: klog.Errorf("resType is not pod or configmap or secret: esType is %s", resType) continue @@ -990,6 +1002,93 @@ func (e *edged) syncPod() { } } +func (e *edged) handleVolume(op string, content []byte) (interface{}, error) { + switch op { + case constants.CSIOperationTypeCreateVolume: + return e.createVolume(content) + case constants.CSIOperationTypeDeleteVolume: + return e.deleteVolume(content) + case constants.CSIOperationTypeControllerPublishVolume: + return e.controllerPublishVolume(content) + case constants.CSIOperationTypeControllerUnpublishVolume: + return e.controllerUnpublishVolume(content) + } + return nil, nil +} + +func (e *edged) createVolume(content []byte) (interface{}, error) { + req := &csi.CreateVolumeRequest{} + err := jsonpb.Unmarshal(bytes.NewReader(content), req) + if err != nil { + klog.Errorf("unmarshal create volume req error: %v", err) + return nil, err + } + + klog.Infof("start create volume: %s", req.Name) + ctl := csiplugin.NewController() + res, err := ctl.CreateVolume(req) + if err != nil { + klog.Errorf("create volume error: %v", err) + return nil, err + } + klog.Infof("end create volume: %s result: %v", req.Name, res) + return res, nil +} + +func (e *edged) deleteVolume(content []byte) (interface{}, error) { + req := &csi.DeleteVolumeRequest{} + err := jsonpb.Unmarshal(bytes.NewReader(content), req) + if err != nil { + klog.Errorf("unmarshal delete volume req error: %v", err) + return nil, err + } + klog.Infof("start delete volume: %s", req.VolumeId) + ctl := csiplugin.NewController() + res, err := ctl.DeleteVolume(req) + if err != nil { + klog.Errorf("delete volume error: %v", err) + return nil, err + } + klog.Infof("end delete volume: %s result: %v", req.VolumeId, res) + return res, nil +} + +func (e *edged) controllerPublishVolume(content []byte) (interface{}, error) { + req := &csi.ControllerPublishVolumeRequest{} + err := jsonpb.Unmarshal(bytes.NewReader(content), req) + if err != nil { + klog.Errorf("unmarshal controller publish volume req error: %v", err) + return nil, err + } + klog.Infof("start controller publish volume: %s", req.VolumeId) + ctl := csiplugin.NewController() + res, err := ctl.ControllerPublishVolume(req) + if err != nil { + klog.Errorf("controller publish volume error: %v", err) + return nil, err + } + klog.Infof("end controller publish volume:: %s result: %v", req.VolumeId, res) + return res, nil +} + +func (e *edged) controllerUnpublishVolume(content []byte) (interface{}, error) { + req := &csi.ControllerUnpublishVolumeRequest{} + err := jsonpb.Unmarshal(bytes.NewReader(content), req) + if err != nil { + klog.Errorf("unmarshal controller publish volume req error: %v", err) + return nil, err + } + klog.Infof("start controller unpublish volume: %s", req.VolumeId) + ctl := csiplugin.NewController() + res, err := ctl.ControllerUnpublishVolume(req) + if err != nil { + klog.Errorf("controller unpublish volume error: %v", err) + return nil, err + } + klog.Infof("end controller unpublish volume:: %s result: %v", req.VolumeId, res) + return res, nil +} + func (e *edged) handlePod(op string, content []byte) (err error) { var pod v1.Pod err = json.Unmarshal(content, &pod) |
