summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authoredisonxiang <xiang.edison@gmail.com>2019-08-21 21:03:28 +0800
committeredisonxiang <xiang.edison@gmail.com>2019-09-05 16:04:36 +0800
commit0d5a632aa8f77b3f3e65fb728d36a112741bc643 (patch)
treeffe68f94c9cb958d971b5a20d39c40dc00cf4fcc
parentMerge pull request #1100 from arcanique/master (diff)
downloadkubeedge-0d5a632aa8f77b3f3e65fb728d36a112741bc643.tar.gz
add csi driver from kubeedge
-rw-r--r--csidriver/cmd/csidriver.go36
-rw-r--r--csidriver/pkg/app/controllerserver.go453
-rw-r--r--csidriver/pkg/app/identityserver.go59
-rw-r--r--csidriver/pkg/app/server.go60
-rw-r--r--csidriver/pkg/app/uds.go65
-rw-r--r--csidriver/pkg/app/utils.go163
6 files changed, 836 insertions, 0 deletions
diff --git a/csidriver/cmd/csidriver.go b/csidriver/cmd/csidriver.go
new file mode 100644
index 000000000..1ea407e15
--- /dev/null
+++ b/csidriver/cmd/csidriver.go
@@ -0,0 +1,36 @@
+package main
+
+import (
+ "flag"
+ "fmt"
+ "os"
+
+ "github.com/kubeedge/kubeedge/csidriver/pkg/app"
+)
+
+func init() {
+ flag.Set("logtostderr", "true")
+}
+
+var (
+ endpoint = flag.String("endpoint", "unix://tmp/csi.sock", "CSI endpoint")
+ driverName = flag.String("drivername", "csidriver", "name of the driver")
+ nodeID = flag.String("nodeid", "", "node id")
+ keEndpoint = flag.String("kubeedgeendpoint", "unix:///kubeedge/kubeedge.sock", "kubeedge endpoint")
+ version = flag.String("version", "", "version")
+)
+
+func main() {
+ flag.Parse()
+ handle()
+ os.Exit(0)
+}
+
+func handle() {
+ driver, err := app.NewCSIDriver(*driverName, *nodeID, *endpoint, *keEndpoint, *version)
+ if err != nil {
+ fmt.Printf("failed to initialize driver: %s", err.Error())
+ os.Exit(1)
+ }
+ driver.Run()
+}
diff --git a/csidriver/pkg/app/controllerserver.go b/csidriver/pkg/app/controllerserver.go
new file mode 100644
index 000000000..9569944c0
--- /dev/null
+++ b/csidriver/pkg/app/controllerserver.go
@@ -0,0 +1,453 @@
+package app
+
+import (
+ "encoding/base64"
+ "encoding/json"
+ "errors"
+ "fmt"
+
+ "github.com/golang/protobuf/jsonpb"
+ "github.com/pborman/uuid"
+ "golang.org/x/net/context"
+ "google.golang.org/grpc/codes"
+ "google.golang.org/grpc/status"
+
+ "github.com/container-storage-interface/spec/lib/go/csi"
+
+ "k8s.io/klog"
+
+ "github.com/kubeedge/beehive/pkg/core/model"
+ "github.com/kubeedge/kubeedge/common/constants"
+)
+
+type controllerServer struct {
+ caps []*csi.ControllerServiceCapability
+ nodeID string
+ keEndpoint string
+}
+
+// NewControllerServer creates controller server
+func NewControllerServer(nodeID, keEndpoint string) *controllerServer {
+ return &controllerServer{
+ caps: getControllerServiceCapabilities(
+ []csi.ControllerServiceCapability_RPC_Type{
+ csi.ControllerServiceCapability_RPC_CREATE_DELETE_VOLUME,
+ csi.ControllerServiceCapability_RPC_PUBLISH_UNPUBLISH_VOLUME,
+ }),
+ nodeID: nodeID,
+ keEndpoint: keEndpoint,
+ }
+}
+
+// CreateVolume issues create volume func
+func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVolumeRequest) (*csi.CreateVolumeResponse, error) {
+ // Check arguments
+ if len(req.GetName()) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "Name missing in request")
+ }
+ caps := req.GetVolumeCapabilities()
+ if caps == nil {
+ return nil, status.Error(codes.InvalidArgument, "Volume Capabilities missing in request")
+ }
+
+ volumeID := uuid.NewUUID().String()
+
+ // Build message struct
+ msg := model.NewMessage("")
+ resource, err := buildResource(cs.nodeID,
+ DefaultNamespace,
+ constants.CSIResourceTypeVolume,
+ volumeID)
+ if err != nil {
+ klog.Errorf("build message resource failed with error: %s", err)
+ return nil, err
+ }
+
+ m := jsonpb.Marshaler{}
+ js, err := m.MarshalToString(req)
+ if err != nil {
+ klog.Errorf("failed to marshal to string with error: %s", err)
+ return nil, err
+ }
+ klog.Infof("create volume marshal to string: %s", js)
+ msg.Content = js
+ msg.BuildRouter(DefaultReceiveModuleName,
+ GroupResource,
+ resource,
+ constants.CSIOperationTypeCreateVolume)
+
+ // Marshal message
+ reqData, err := json.Marshal(msg)
+ if err != nil {
+ klog.Errorf("marshal request failed with error: %v", err)
+ return nil, err
+ }
+
+ // Send message to KubeEdge
+ resdata, err := send2KubeEdge(string(reqData), cs.keEndpoint)
+ if err != nil {
+ klog.Errorf("send to kubeedge failed with error: %v", err)
+ return nil, err
+ }
+
+ // Unmarshal message
+ result, err := extractMessage(resdata)
+ if err != nil {
+ klog.Errorf("unmarshal response failed with error: %v", err)
+ return nil, err
+ }
+
+ klog.Infof("create volume result: %v", result)
+ data := result.GetContent().(string)
+
+ if msg.GetOperation() == model.ResponseErrorOperation {
+ klog.Errorf("create volume with error: %s", data)
+ return nil, errors.New(data)
+ }
+
+ decodeBytes, err := base64.StdEncoding.DecodeString(data)
+ if err != nil {
+ klog.Errorf("create volume decode with error: %v", err)
+ return nil, err
+ }
+
+ response := &csi.CreateVolumeResponse{}
+ err = json.Unmarshal([]byte(decodeBytes), response)
+ if err != nil {
+ klog.Errorf("create volume unmarshal with error: %v", err)
+ return nil, nil
+ }
+ klog.Infof("create volume response: %v", response)
+
+ createVolumeResponse := &csi.CreateVolumeResponse{}
+ if req.GetVolumeContentSource() != nil {
+ createVolumeResponse = &csi.CreateVolumeResponse{
+ Volume: &csi.Volume{
+ VolumeId: response.Volume.VolumeId,
+ CapacityBytes: req.GetCapacityRange().GetRequiredBytes(),
+ VolumeContext: req.GetParameters(),
+ ContentSource: req.GetVolumeContentSource(),
+ },
+ }
+ } else {
+ createVolumeResponse = &csi.CreateVolumeResponse{
+ Volume: &csi.Volume{
+ VolumeId: response.Volume.VolumeId,
+ CapacityBytes: req.GetCapacityRange().GetRequiredBytes(),
+ VolumeContext: req.GetParameters(),
+ },
+ }
+ }
+ return createVolumeResponse, nil
+}
+
+// DeleteVolume issues delete volume func
+func (cs *controllerServer) DeleteVolume(ctx context.Context, req *csi.DeleteVolumeRequest) (*csi.DeleteVolumeResponse, error) {
+ // Check arguments
+ if len(req.GetVolumeId()) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "Volume ID missing in request")
+ }
+
+ // Build message struct
+ msg := model.NewMessage("")
+ resource, err := buildResource(cs.nodeID,
+ DefaultNamespace,
+ constants.CSIResourceTypeVolume,
+ req.GetVolumeId())
+ if err != nil {
+ klog.Errorf("build message resource failed with error: %s", err)
+ return nil, err
+ }
+
+ m := jsonpb.Marshaler{}
+ js, err := m.MarshalToString(req)
+ if err != nil {
+ klog.Errorf("failed to marshal to string with error: %s", err)
+ return nil, err
+ }
+ klog.Infof("delete volume marshal to string: %s", js)
+ msg.Content = js
+ msg.BuildRouter(DefaultReceiveModuleName,
+ GroupResource,
+ resource,
+ constants.CSIOperationTypeDeleteVolume)
+
+ // Marshal message
+ reqData, err := json.Marshal(msg)
+ if err != nil {
+ klog.Errorf("marshal request failed with error: %v", err)
+ return nil, err
+ }
+
+ // Send message to KubeEdge
+ resdata, err := send2KubeEdge(string(reqData), cs.keEndpoint)
+ if err != nil {
+ klog.Errorf("send to kubeedge failed with error: %v", err)
+ return nil, err
+ }
+
+ // Unmarshal message
+ result, err := extractMessage(resdata)
+ if err != nil {
+ klog.Errorf("unmarshal response failed with error: %v", err)
+ return nil, err
+ }
+
+ klog.Infof("delete volume result: %v", result)
+ data := result.GetContent().(string)
+
+ if msg.GetOperation() == model.ResponseErrorOperation {
+ klog.Errorf("delete volume with error: %s", data)
+ return nil, errors.New(data)
+ }
+
+ decodeBytes, err := base64.StdEncoding.DecodeString(data)
+ if err != nil {
+ klog.Errorf("delete volume decode with error: %v", err)
+ return nil, err
+ }
+
+ deleteVolumeResponse := &csi.DeleteVolumeResponse{}
+ err = json.Unmarshal([]byte(decodeBytes), deleteVolumeResponse)
+ if err != nil {
+ klog.Errorf("delete volume unmarshal with error: %v", err)
+ return nil, nil
+ }
+ klog.Infof("delete volume response: %v", deleteVolumeResponse)
+ return deleteVolumeResponse, nil
+}
+
+// ControllerPublishVolume issues controller publish volume func
+func (cs *controllerServer) ControllerPublishVolume(ctx context.Context, req *csi.ControllerPublishVolumeRequest) (*csi.ControllerPublishVolumeResponse, error) {
+ instanceID := req.GetNodeId()
+ volumeID := req.GetVolumeId()
+
+ if len(volumeID) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "ControllerPublishVolume Volume ID must be provided")
+ }
+
+ if len(instanceID) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "ControllerPublishVolume Instance ID must be provided")
+ }
+
+ // Build message struct
+ msg := model.NewMessage("")
+ resource, err := buildResource(cs.nodeID,
+ DefaultNamespace,
+ constants.CSIResourceTypeVolume,
+ volumeID)
+ if err != nil {
+ klog.Errorf("build message resource failed with error: %s", err)
+ return nil, err
+ }
+
+ m := jsonpb.Marshaler{}
+ js, err := m.MarshalToString(req)
+ if err != nil {
+ klog.Errorf("failed to marshal to string with error: %s", err)
+ return nil, err
+ }
+ klog.Infof("controller publish volume marshal to string: %s", js)
+ msg.Content = js
+ msg.BuildRouter(DefaultReceiveModuleName,
+ GroupResource,
+ resource,
+ constants.CSIOperationTypeControllerPublishVolume)
+
+ // Marshal message
+ reqData, err := json.Marshal(msg)
+ if err != nil {
+ klog.Errorf("marshal request failed with error: %v", err)
+ return nil, err
+ }
+
+ // Send message to KubeEdge
+ resdata, err := send2KubeEdge(string(reqData), cs.keEndpoint)
+ if err != nil {
+ klog.Errorf("send to kubeedge failed with error: %v", err)
+ return nil, err
+ }
+
+ // Unmarshal message
+ result, err := extractMessage(resdata)
+ if err != nil {
+ klog.Errorf("unmarshal response failed with error: %v", err)
+ return nil, err
+ }
+
+ klog.Infof("controller publish volume result: %v", result)
+ data := result.GetContent().(string)
+
+ if msg.GetOperation() == model.ResponseErrorOperation {
+ klog.Errorf("controller publish volume with error: %s", data)
+ return nil, errors.New(data)
+ }
+
+ decodeBytes, err := base64.StdEncoding.DecodeString(data)
+ if err != nil {
+ klog.Errorf("controller publish volume decode with error: %v", err)
+ return nil, err
+ }
+
+ controllerPublishVolumeResponse := &csi.ControllerPublishVolumeResponse{}
+ err = json.Unmarshal([]byte(decodeBytes), controllerPublishVolumeResponse)
+ if err != nil {
+ klog.Errorf("controller publish volume unmarshal with error: %v", err)
+ return nil, nil
+ }
+ klog.Infof("controller publish volume response: %v", controllerPublishVolumeResponse)
+ return controllerPublishVolumeResponse, nil
+}
+
+// ControllerUnpublishVolume issues controller unpublish volume func
+func (cs *controllerServer) ControllerUnpublishVolume(ctx context.Context, req *csi.ControllerUnpublishVolumeRequest) (*csi.ControllerUnpublishVolumeResponse, error) {
+ instanceID := req.GetNodeId()
+ volumeID := req.GetVolumeId()
+
+ if len(volumeID) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "ControllerUnpublishVolume Volume ID must be provided")
+ }
+
+ if len(instanceID) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "ControllerUnpublishVolume Instance ID must be provided")
+ }
+
+ // Build message struct
+ msg := model.NewMessage("")
+ resource, err := buildResource(cs.nodeID,
+ DefaultNamespace,
+ constants.CSIResourceTypeVolume,
+ volumeID)
+ if err != nil {
+ klog.Errorf("Build message resource failed with error: %s", err)
+ return nil, err
+ }
+
+ m := jsonpb.Marshaler{}
+ js, err := m.MarshalToString(req)
+ if err != nil {
+ klog.Errorf("failed to marshal to string with error: %s", err)
+ return nil, err
+ }
+ klog.Infof("controller Unpublish Volume marshal to string: %s", js)
+ msg.Content = js
+ msg.BuildRouter(DefaultReceiveModuleName,
+ GroupResource,
+ resource,
+ constants.CSIOperationTypeControllerUnpublishVolume)
+
+ // Marshal message
+ reqData, err := json.Marshal(msg)
+ if err != nil {
+ klog.Errorf("marshal request failed with error: %v", err)
+ return nil, err
+ }
+
+ // Send message to KubeEdge
+ resdata, err := send2KubeEdge(string(reqData), cs.keEndpoint)
+ if err != nil {
+ klog.Errorf("send to kubeedge failed with error: %v", err)
+ return nil, err
+ }
+
+ // Unmarshal message
+ result, err := extractMessage(resdata)
+ if err != nil {
+ klog.Errorf("unmarshal response failed with error: %v", err)
+ return nil, err
+ }
+
+ klog.Infof("controller Unpublish Volume result: %v", result)
+ data := result.GetContent().(string)
+
+ if msg.GetOperation() == model.ResponseErrorOperation {
+ klog.Errorf("controller Unpublish Volume with error: %s", data)
+ return nil, errors.New(data)
+ }
+
+ decodeBytes, err := base64.StdEncoding.DecodeString(data)
+ if err != nil {
+ klog.Errorf("controller Unpublish Volume decode with error: %v", err)
+ return nil, err
+ }
+
+ controllerUnpublishVolumeResponse := &csi.ControllerUnpublishVolumeResponse{}
+ err = json.Unmarshal([]byte(decodeBytes), controllerUnpublishVolumeResponse)
+ if err != nil {
+ klog.Errorf("controller Unpublish Volume unmarshal with error: %v", err)
+ return nil, nil
+ }
+ klog.Infof("controller Unpublish Volume response: %v", controllerUnpublishVolumeResponse)
+ return controllerUnpublishVolumeResponse, nil
+}
+
+func (cs *controllerServer) ValidateVolumeCapabilities(ctx context.Context, req *csi.ValidateVolumeCapabilitiesRequest) (*csi.ValidateVolumeCapabilitiesResponse, error) {
+ // Check arguments
+ if len(req.GetVolumeId()) == 0 {
+ return nil, status.Error(codes.InvalidArgument, "Volume ID cannot be empty")
+ }
+ if len(req.VolumeCapabilities) == 0 {
+ return nil, status.Error(codes.InvalidArgument, req.VolumeId)
+ }
+
+ for _, cap := range req.GetVolumeCapabilities() {
+ if cap.GetMount() == nil && cap.GetBlock() == nil {
+ return nil, status.Error(codes.InvalidArgument, "cannot have both mount and block access type be undefined")
+ }
+ }
+
+ return &csi.ValidateVolumeCapabilitiesResponse{
+ Confirmed: &csi.ValidateVolumeCapabilitiesResponse_Confirmed{
+ VolumeContext: req.GetVolumeContext(),
+ VolumeCapabilities: req.GetVolumeCapabilities(),
+ Parameters: req.GetParameters(),
+ },
+ }, nil
+}
+
+func (cs *controllerServer) ControllerGetCapabilities(ctx context.Context, req *csi.ControllerGetCapabilitiesRequest) (*csi.ControllerGetCapabilitiesResponse, error) {
+ return &csi.ControllerGetCapabilitiesResponse{
+ Capabilities: cs.caps,
+ }, nil
+}
+
+func getControllerServiceCapabilities(cl []csi.ControllerServiceCapability_RPC_Type) []*csi.ControllerServiceCapability {
+ var csc []*csi.ControllerServiceCapability
+
+ for _, cap := range cl {
+ klog.Infof("Enabling controller service capability: %v", cap.String())
+ csc = append(csc, &csi.ControllerServiceCapability{
+ Type: &csi.ControllerServiceCapability_Rpc{
+ Rpc: &csi.ControllerServiceCapability_RPC{
+ Type: cap,
+ },
+ },
+ })
+ }
+
+ return csc
+}
+
+func (cs *controllerServer) GetCapacity(ctx context.Context, req *csi.GetCapacityRequest) (*csi.GetCapacityResponse, error) {
+ return nil, status.Error(codes.Unimplemented, "")
+}
+
+func (cs *controllerServer) ListVolumes(ctx context.Context, req *csi.ListVolumesRequest) (*csi.ListVolumesResponse, error) {
+ return nil, status.Error(codes.Unimplemented, "")
+}
+
+func (cs *controllerServer) ControllerExpandVolume(ctx context.Context, req *csi.ControllerExpandVolumeRequest) (*csi.ControllerExpandVolumeResponse, error) {
+ return nil, status.Error(codes.Unimplemented, fmt.Sprintf("ControllerExpandVolume is not yet implemented"))
+}
+
+func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateSnapshotRequest) (*csi.CreateSnapshotResponse, error) {
+ return nil, status.Error(codes.Unimplemented, fmt.Sprintf("CreateSnapshot is not yet implemented"))
+}
+
+func (cs *controllerServer) DeleteSnapshot(ctx context.Context, req *csi.DeleteSnapshotRequest) (*csi.DeleteSnapshotResponse, error) {
+ return nil, status.Error(codes.Unimplemented, fmt.Sprintf("DeleteSnapshot is not yet implemented"))
+}
+
+func (cs *controllerServer) ListSnapshots(ctx context.Context, req *csi.ListSnapshotsRequest) (*csi.ListSnapshotsResponse, error) {
+ return nil, status.Error(codes.Unimplemented, fmt.Sprintf("ListSnapshots is not yet implemented"))
+}
diff --git a/csidriver/pkg/app/identityserver.go b/csidriver/pkg/app/identityserver.go
new file mode 100644
index 000000000..11d455bda
--- /dev/null
+++ b/csidriver/pkg/app/identityserver.go
@@ -0,0 +1,59 @@
+package app
+
+import (
+ "github.com/container-storage-interface/spec/lib/go/csi"
+ "golang.org/x/net/context"
+ "google.golang.org/grpc/codes"
+ "google.golang.org/grpc/status"
+
+ "k8s.io/klog"
+)
+
+type identityServer struct {
+ name string
+ version string
+}
+
+// NewIdentityServer creates identity server
+func NewIdentityServer(name, version string) *identityServer {
+ return &identityServer{
+ name: name,
+ version: version,
+ }
+}
+
+func (ids *identityServer) GetPluginInfo(ctx context.Context, req *csi.GetPluginInfoRequest) (*csi.GetPluginInfoResponse, error) {
+ klog.Info("using default GetPluginInfo")
+
+ if ids.name == "" {
+ return nil, status.Error(codes.Unavailable, "driver name not configured")
+ }
+
+ if ids.version == "" {
+ return nil, status.Error(codes.Unavailable, "driver is missing version")
+ }
+
+ return &csi.GetPluginInfoResponse{
+ Name: ids.name,
+ VendorVersion: ids.version,
+ }, nil
+}
+
+func (ids *identityServer) Probe(ctx context.Context, req *csi.ProbeRequest) (*csi.ProbeResponse, error) {
+ return &csi.ProbeResponse{}, nil
+}
+
+func (ids *identityServer) GetPluginCapabilities(ctx context.Context, req *csi.GetPluginCapabilitiesRequest) (*csi.GetPluginCapabilitiesResponse, error) {
+ klog.Info("using default capabilities")
+ return &csi.GetPluginCapabilitiesResponse{
+ Capabilities: []*csi.PluginCapability{
+ {
+ Type: &csi.PluginCapability_Service_{
+ Service: &csi.PluginCapability_Service{
+ Type: csi.PluginCapability_Service_CONTROLLER_SERVICE,
+ },
+ },
+ },
+ },
+ }, nil
+}
diff --git a/csidriver/pkg/app/server.go b/csidriver/pkg/app/server.go
new file mode 100644
index 000000000..edb2007cd
--- /dev/null
+++ b/csidriver/pkg/app/server.go
@@ -0,0 +1,60 @@
+package app
+
+import (
+ "fmt"
+
+ "k8s.io/klog"
+)
+
+type CSIDriver struct {
+ name string
+ nodeID string
+ version string
+ endpoint string
+ keEndpoint string
+
+ ids *identityServer
+ cs *controllerServer
+}
+
+var (
+ vendorVersion = "dev"
+)
+
+// NewCSIDriver creates a new server
+func NewCSIDriver(driverName, nodeID, endpoint, keEndpoint, version string) (*CSIDriver, error) {
+ if driverName == "" {
+ return nil, fmt.Errorf("no driver name provided")
+ }
+ if nodeID == "" {
+ return nil, fmt.Errorf("no node id provided")
+ }
+ if endpoint == "" {
+ return nil, fmt.Errorf("no driver endpoint provided")
+ }
+ if keEndpoint == "" {
+ return nil, fmt.Errorf("no kubeedge endpoint provided")
+ }
+ if version != "" {
+ vendorVersion = version
+ }
+ klog.Infof("driver: %s version: %s", driverName, vendorVersion)
+
+ return &CSIDriver{
+ name: driverName,
+ version: vendorVersion,
+ nodeID: nodeID,
+ endpoint: endpoint,
+ keEndpoint: keEndpoint,
+ }, nil
+}
+
+func (cd *CSIDriver) Run() {
+ // Create GRPC servers
+ cd.ids = NewIdentityServer(cd.name, cd.version)
+ cd.cs = NewControllerServer(cd.nodeID, cd.keEndpoint)
+
+ s := NewNonBlockingGRPCServer()
+ s.Start(cd.endpoint, cd.ids, cd.cs, nil)
+ s.Wait()
+}
diff --git a/csidriver/pkg/app/uds.go b/csidriver/pkg/app/uds.go
new file mode 100644
index 000000000..48d0c19e1
--- /dev/null
+++ b/csidriver/pkg/app/uds.go
@@ -0,0 +1,65 @@
+package app
+
+import (
+ "net"
+
+ "k8s.io/klog"
+)
+
+const (
+ // DefaultBufferSize represents default buffer size
+ DefaultBufferSize = 10480
+)
+
+// UnixDomainSocket struct
+type UnixDomainSocket struct {
+ filename string
+ buffersize int
+}
+
+// NewUnixDomainSocket create new socket
+func NewUnixDomainSocket(filename string, buffersize ...int) *UnixDomainSocket {
+ size := DefaultBufferSize
+ if buffersize != nil {
+ size = buffersize[0]
+ }
+ us := UnixDomainSocket{filename: filename, buffersize: size}
+ return &us
+}
+
+// Connect for client
+func (us *UnixDomainSocket) Connect() (net.Conn, error) {
+ // parse
+ proto, addr, err := parseEndpoint(us.filename)
+ if err != nil {
+ klog.Errorf("failed to parseEndpoint: %v", err)
+ return nil, err
+ }
+
+ // dial
+ c, err := net.Dial(proto, addr)
+ if err != nil {
+ klog.Errorf("failed to dial: %v", err)
+ return nil, err
+ }
+ return c, nil
+}
+
+// Send msg for client
+func (us *UnixDomainSocket) Send(c net.Conn, context string) (string, error) {
+ // send msg
+ _, err := c.Write([]byte(context))
+ if err != nil {
+ klog.Errorf("failed to write buffer: %v", err)
+ return "", err
+ }
+
+ // read response
+ buf := make([]byte, us.buffersize)
+ nr, err := c.Read(buf)
+ if err != nil {
+ klog.Errorf("failed to read buffer: %v", err)
+ return "", err
+ }
+ return string(buf[0:nr]), nil
+}
diff --git a/csidriver/pkg/app/utils.go b/csidriver/pkg/app/utils.go
new file mode 100644
index 000000000..b87e397a3
--- /dev/null
+++ b/csidriver/pkg/app/utils.go
@@ -0,0 +1,163 @@
+package app
+
+import (
+ "encoding/json"
+ "errors"
+ "fmt"
+ "net"
+ "os"
+ "strings"
+ "sync"
+
+ "k8s.io/klog"
+
+ "golang.org/x/net/context"
+ "google.golang.org/grpc"
+
+ "github.com/container-storage-interface/spec/lib/go/csi"
+ "github.com/kubernetes-csi/csi-lib-utils/protosanitizer"
+
+ "github.com/kubeedge/beehive/pkg/core/model"
+ "github.com/kubeedge/kubeedge/common/constants"
+)
+
+// Constant defines csi related parameters
+const (
+ GroupResource = "resource"
+ DefaultNamespace = "default"
+ DefaultReceiveModuleName = "cloudhub"
+)
+
+// NewNonBlockingGRPCServer creates a new nonblocking server
+func NewNonBlockingGRPCServer() *nonBlockingGRPCServer {
+ return &nonBlockingGRPCServer{}
+}
+
+// NonBlocking server
+type nonBlockingGRPCServer struct {
+ wg sync.WaitGroup
+ server *grpc.Server
+}
+
+func (s *nonBlockingGRPCServer) Start(endpoint string, ids csi.IdentityServer, cs csi.ControllerServer, ns csi.NodeServer) {
+ s.wg.Add(1)
+ go s.serve(endpoint, ids, cs, ns)
+ return
+}
+
+func (s *nonBlockingGRPCServer) Wait() {
+ s.wg.Wait()
+}
+
+func (s *nonBlockingGRPCServer) Stop() {
+ s.server.GracefulStop()
+}
+
+func (s *nonBlockingGRPCServer) ForceStop() {
+ s.server.Stop()
+}
+
+func (s *nonBlockingGRPCServer) serve(endpoint string, ids csi.IdentityServer, cs csi.ControllerServer, ns csi.NodeServer) {
+
+ proto, addr, err := parseEndpoint(endpoint)
+ if err != nil {
+ klog.Fatal(err.Error())
+ }
+
+ if proto == "unix" {
+ addr = "/" + addr
+ if err := os.Remove(addr); err != nil && !os.IsNotExist(err) {
+ klog.Errorf("failed to remove %s, error: %s", addr, err.Error())
+ }
+ }
+
+ listener, err := net.Listen(proto, addr)
+ if err != nil {
+ klog.Errorf("failed to listen: %v", err)
+ }
+
+ opts := []grpc.ServerOption{
+ grpc.UnaryInterceptor(logGRPC),
+ }
+ server := grpc.NewServer(opts...)
+ s.server = server
+
+ if ids != nil {
+ csi.RegisterIdentityServer(server, ids)
+ }
+ if cs != nil {
+ csi.RegisterControllerServer(server, cs)
+ }
+ if ns != nil {
+ csi.RegisterNodeServer(server, ns)
+ }
+ klog.Infof("listening for connections on address: %#v", listener.Addr())
+ server.Serve(listener)
+}
+
+func parseEndpoint(ep string) (string, string, error) {
+ if strings.HasPrefix(strings.ToLower(ep), "unix://") || strings.HasPrefix(strings.ToLower(ep), "tcp://") {
+ s := strings.SplitN(ep, "://", 2)
+ if s[1] != "" {
+ return s[0], s[1], nil
+ }
+ }
+ return "", "", fmt.Errorf("invalid endpoint: %v", ep)
+}
+
+func logGRPC(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
+ klog.Infof("gprc call: %s", info.FullMethod)
+ klog.Infof("gprc request: %+v", protosanitizer.StripSecrets(req))
+ resp, err := handler(ctx, req)
+ if err != nil {
+ klog.Errorf("gprc error: %v", err)
+ } else {
+ klog.Infof("gprc response: %+v", protosanitizer.StripSecrets(resp))
+ }
+ return resp, err
+}
+
+// buildResource return a string as "beehive/pkg/core/model".Message.Router.Resource
+func buildResource(nodeID, namespace, resourceType, resourceID string) (resource string, err error) {
+ if nodeID == "" || namespace == "" || resourceType == "" {
+ err = fmt.Errorf("required parameter are not set (node id, namespace or resource type)")
+ return
+ }
+ resource = fmt.Sprintf("%s%s%s%s%s%s%s", "node", constants.ResourceSep, nodeID, constants.ResourceSep, namespace, constants.ResourceSep, resourceType)
+ if resourceID != "" {
+ resource += fmt.Sprintf("%s%s", constants.ResourceSep, resourceID)
+ }
+ return
+}
+
+// send2KubeEdge sends messages to KubeEdge
+func send2KubeEdge(context, keEndpoint string) (string, error) {
+ us := NewUnixDomainSocket(keEndpoint)
+ // connect
+ r, err := us.Connect()
+ if err != nil {
+ return "", err
+ }
+ // send
+ res, err := us.Send(r, context)
+ if err != nil {
+ return "", err
+ }
+ return res, nil
+}
+
+// extractMessage extracts message
+func extractMessage(context string) (*model.Message, error) {
+ var msg *model.Message
+ if context != "" {
+ err := json.Unmarshal([]byte(context), &msg)
+ if err != nil {
+ return nil, err
+ }
+ } else {
+ err := errors.New("failed to extract message with empty context")
+ klog.Errorf("%v", err)
+ return nil, err
+ }
+ return msg, nil
+}