summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authoredisonxiang <xiang.edison@gmail.com>2019-08-22 00:53:38 +0800
committeredisonxiang <xiang.edison@gmail.com>2019-09-04 10:09:15 +0800
commit8939add58f4d560f55ed2b126483a04e355c32ff (patch)
tree8ea9c970c508675460dbc64f130556cf6a5b1619 /edge
parentupdate required vendors (diff)
downloadkubeedge-8939add58f4d560f55ed2b126483a04e355c32ff.tar.gz
invoke func from edged
Diffstat (limited to 'edge')
-rw-r--r--edge/pkg/edged/edged.go99
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)