diff options
| author | edisonxiang <xiang.edison@gmail.com> | 2019-08-21 21:03:28 +0800 |
|---|---|---|
| committer | edisonxiang <xiang.edison@gmail.com> | 2019-09-05 16:04:36 +0800 |
| commit | 0d5a632aa8f77b3f3e65fb728d36a112741bc643 (patch) | |
| tree | ffe68f94c9cb958d971b5a20d39c40dc00cf4fcc | |
| parent | Merge pull request #1100 from arcanique/master (diff) | |
| download | kubeedge-0d5a632aa8f77b3f3e65fb728d36a112741bc643.tar.gz | |
add csi driver from kubeedge
| -rw-r--r-- | csidriver/cmd/csidriver.go | 36 | ||||
| -rw-r--r-- | csidriver/pkg/app/controllerserver.go | 453 | ||||
| -rw-r--r-- | csidriver/pkg/app/identityserver.go | 59 | ||||
| -rw-r--r-- | csidriver/pkg/app/server.go | 60 | ||||
| -rw-r--r-- | csidriver/pkg/app/uds.go | 65 | ||||
| -rw-r--r-- | csidriver/pkg/app/utils.go | 163 |
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 +} |
