summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorvincentgoat <linguohui1@huawei.com>2022-06-25 17:21:07 +0800
committervincentgoat <linguohui1@huawei.com>2022-07-11 21:15:05 +0800
commit412471db1bff0ace3fd2ece437eecce2b4811cee (patch)
treecd4df7498eeb73298a3ff8f4e769b04d35509e50
parentMerge pull request #4007 from gy95/automated-cherry-pick-of-#3956-upstream-re... (diff)
downloadkubeedge-412471db1bff0ace3fd2ece437eecce2b4811cee.tar.gz
fix complex parameters
Signed-off-by: vincentgoat <linguohui1@huawei.com>
-rw-r--r--edge/pkg/devicetwin/dtmanager/communicate.go13
-rw-r--r--edge/pkg/devicetwin/dtmanager/communicate_test.go2
-rw-r--r--edge/pkg/devicetwin/dtmanager/membership.go10
-rw-r--r--edge/pkg/devicetwin/dtmanager/membership_test.go12
-rw-r--r--edge/pkg/devicetwin/dtmanager/twin.go17
-rw-r--r--edge/pkg/devicetwin/dtmanager/twin_test.go110
-rw-r--r--edge/pkg/edged/edged_pods.go4
-rw-r--r--edge/pkg/edged/edged_status.go9
-rw-r--r--edge/pkg/edged/volume/csi/csi_attacher.go7
-rw-r--r--edge/pkg/edged/volume/csi/csi_plugin.go61
-rw-r--r--edgesite/cmd/edgesite-server/app/server.go19
-rw-r--r--keadm/cmd/keadm/app/cmd/util/common.go2
-rw-r--r--keadm/cmd/keadm/app/cmd/util/debinstaller.go2
-rw-r--r--keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go2
-rw-r--r--keadm/cmd/keadm/app/cmd/util/rpminstaller.go2
-rw-r--r--staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go23
-rw-r--r--staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go21
17 files changed, 153 insertions, 163 deletions
diff --git a/edge/pkg/devicetwin/dtmanager/communicate.go b/edge/pkg/devicetwin/dtmanager/communicate.go
index 787a8d178..d8d459d74 100644
--- a/edge/pkg/devicetwin/dtmanager/communicate.go
+++ b/edge/pkg/devicetwin/dtmanager/communicate.go
@@ -53,7 +53,7 @@ func (cw CommWorker) Start() {
}
case <-time.After(time.Duration(60) * time.Second):
- cw.checkConfirm(cw.DTContexts, nil)
+ cw.checkConfirm(cw.DTContexts)
case v, ok := <-cw.HeartBeatChan:
if !ok {
return
@@ -100,7 +100,7 @@ func dealLifeCycle(context *dtcontext.DTContext, resource string, msg interface{
connectedInfo, _ := message.Content.(string)
if strings.Compare(connectedInfo, connect.CloudConnected) == 0 {
if strings.Compare(context.State, dtcommon.Disconnected) == 0 {
- _, err := detailRequest(context, msg)
+ err := detailRequest(context)
if err != nil {
klog.Errorf("detail request: %v", err)
return err
@@ -126,7 +126,7 @@ func dealConfirm(context *dtcontext.DTContext, resource string, msg interface{})
return nil
}
-func detailRequest(context *dtcontext.DTContext, msg interface{}) (interface{}, error) {
+func detailRequest(context *dtcontext.DTContext) error {
getDetail := dttype.GetDetailNode{
EventType: "group_membership_event",
EventID: uuid.New().String(),
@@ -136,7 +136,7 @@ func detailRequest(context *dtcontext.DTContext, msg interface{}) (interface{},
getDetailJSON, marshalErr := json.Marshal(getDetail)
if marshalErr != nil {
klog.Errorf("Marshal request error while request detail, err: %#v", marshalErr)
- return nil, marshalErr
+ return marshalErr
}
message := context.BuildModelMessage("resource", "", "membership/detail", "get", string(getDetailJSON))
@@ -144,10 +144,10 @@ func detailRequest(context *dtcontext.DTContext, msg interface{}) (interface{},
msgID := message.GetID()
context.ConfirmMap.Store(msgID, &dttype.DTMessage{Msg: message, Action: dtcommon.SendToCloud, Type: dtcommon.CommModule})
beehiveContext.Send(dtcommon.HubModule, *message)
- return nil, nil
+ return nil
}
-func (cw CommWorker) checkConfirm(context *dtcontext.DTContext, msg interface{}) (interface{}, error) {
+func (cw CommWorker) checkConfirm(context *dtcontext.DTContext) {
klog.V(2).Info("CheckConfirm")
context.ConfirmMap.Range(func(key interface{}, value interface{}) bool {
dtmsg, ok := value.(*dttype.DTMessage)
@@ -166,5 +166,4 @@ func (cw CommWorker) checkConfirm(context *dtcontext.DTContext, msg interface{})
}
return true
})
- return nil, nil
}
diff --git a/edge/pkg/devicetwin/dtmanager/communicate_test.go b/edge/pkg/devicetwin/dtmanager/communicate_test.go
index 022605f75..99ea3f2be 100644
--- a/edge/pkg/devicetwin/dtmanager/communicate_test.go
+++ b/edge/pkg/devicetwin/dtmanager/communicate_test.go
@@ -351,7 +351,7 @@ func TestCheckConfirm(t *testing.T) {
cw := CommWorker{
Worker: test.Worker,
}
- cw.checkConfirm(test.context, test.msg)
+ cw.checkConfirm(test.context)
_, exist := test.context.ConfirmMap.Load("actionMessage")
if !exist {
t.Errorf(" checkconfirm() failed because dealSendToCloud() failed to store message in ConfirmMap")
diff --git a/edge/pkg/devicetwin/dtmanager/membership.go b/edge/pkg/devicetwin/dtmanager/membership.go
index 792b91b33..0d1df0b17 100644
--- a/edge/pkg/devicetwin/dtmanager/membership.go
+++ b/edge/pkg/devicetwin/dtmanager/membership.go
@@ -158,7 +158,9 @@ func dealMembershipGet(context *dtcontext.DTContext, resource string, msg interf
return errors.New("assertion failed")
}
- dealMembershipGetInner(context, contentData)
+ if err := dealMembershipGetInner(context, contentData); err != nil {
+ return err
+ }
return nil
}
@@ -345,11 +347,13 @@ func dealMembershipGetInner(context *dtcontext.DTContext, payload []byte) error
topic := dtcommon.MemETPrefix + context.NodeName + dtcommon.MemETGetResultSuffix
klog.Infof("Deal getting membership successful and send the result")
- context.Send("",
+ err = context.Send("",
dtcommon.SendToEdge,
dtcommon.CommModule,
context.BuildModelMessage(modules.BusGroup, "", topic, messagepkg.OperationPublish, result))
-
+ if err != nil {
+ return err
+ }
return nil
}
diff --git a/edge/pkg/devicetwin/dtmanager/membership_test.go b/edge/pkg/devicetwin/dtmanager/membership_test.go
index 62ab3225b..944afb339 100644
--- a/edge/pkg/devicetwin/dtmanager/membership_test.go
+++ b/edge/pkg/devicetwin/dtmanager/membership_test.go
@@ -301,8 +301,8 @@ func TestDealMembershipGetValid(t *testing.T) {
Content: content,
}
err := dealMembershipGet(dtc, "t", m)
- if err != nil {
- t.Errorf("expected nil, but got error: %v", err)
+ if !reflect.DeepEqual(err, errors.New("Not found chan to communicate")) {
+ t.Errorf("expected %v, but got error: %v", errors.New("Not found chan to communicate"), err)
}
}
@@ -329,8 +329,8 @@ func TestDealMembershipGetInnerValid(t *testing.T) {
content, _ := json.Marshal(payload)
err := dealMembershipGetInner(dtc, content)
- if err != nil {
- t.Errorf("expected nil, but got error: %v", err)
+ if !reflect.DeepEqual(err, errors.New("Not found chan to communicate")) {
+ t.Errorf("expected %v, but got error: %v", errors.New("Not found chan to communicate"), err)
}
}
@@ -343,8 +343,8 @@ func TestDealMembershipGetInnerInValid(t *testing.T) {
}
err := dealMembershipGetInner(dtc, []byte("invalid"))
- if err != nil {
- t.Errorf("expected nil, but got error: %v", err)
+ if !reflect.DeepEqual(err, errors.New("Not found chan to communicate")) {
+ t.Errorf("expected %v, but got error: %v", errors.New("Not found chan to communicate"), err)
}
}
diff --git a/edge/pkg/devicetwin/dtmanager/twin.go b/edge/pkg/devicetwin/dtmanager/twin.go
index 48faa94d3..d84728fc3 100644
--- a/edge/pkg/devicetwin/dtmanager/twin.go
+++ b/edge/pkg/devicetwin/dtmanager/twin.go
@@ -22,7 +22,7 @@ import (
const (
//RestDealType update from mqtt
RestDealType = 0
- //SyncDealType update form cloud sync
+ //SyncDealType update from cloud sync
SyncDealType = 1
//DetailDealType detail update from cloud
DetailDealType = 2
@@ -367,7 +367,7 @@ func dealVersion(version *dttype.TwinVersion, reqVersion *dttype.TwinVersion, de
return true, nil
}
-func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key string, twin *dttype.MsgTwin, msgTwin *dttype.MsgTwin, dealType int) error {
+func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key string, twin *dttype.MsgTwin, msgTwin *dttype.MsgTwin, dealType int) {
document := returnResult.Document
document[key] = &dttype.TwinDoc{}
copytwin := dttype.CopyMsgTwin(twin, true)
@@ -379,7 +379,7 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st
isChange := false
if msgTwin == nil && dealType == RestDealType && *twin.Optional || dealType >= SyncDealType && strings.Compare(msgTwin.Metadata.Type, dtcommon.TypeDeleted) == 0 {
if twin.Metadata != nil && strings.Compare(twin.Metadata.Type, dtcommon.TypeDeleted) == 0 {
- return nil
+ return
}
if dealType != RestDealType {
dealType = SyncTwinDeleteDealType
@@ -403,7 +403,7 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st
syncResult[key] = &copySync
delete(document, key)
returnResult.SyncResult = syncResult
- return nil
+ return
}
} else {
expectedVersionJSON, _ := json.Marshal(expectedVersion)
@@ -445,7 +445,7 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st
syncResult[key] = &copySync
delete(document, key)
returnResult.SyncResult = syncResult
- return nil
+ return
}
} else {
actualVersionJSON, _ := json.Marshal(actualVersion)
@@ -485,8 +485,6 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st
delete(document, key)
delete(syncResult, key)
}
-
- return nil
}
//0:expected ,1 :actual
@@ -945,10 +943,7 @@ func DealMsgTwin(context *dtcontext.DTContext, deviceID string, msgTwins map[str
klog.Infof("Not found metadata of twin")
}
if msgTwin == nil && dealType == 0 || dealType >= 1 && strings.Compare(msgTwin.Metadata.Type, dtcommon.TypeDeleted) == 0 {
- err = dealTwinDelete(&returnResult, deviceID, key, twin, msgTwin, dealType)
- if err != nil {
- return returnResult
- }
+ dealTwinDelete(&returnResult, deviceID, key, twin, msgTwin, dealType)
continue
}
err = dealTwinCompare(&returnResult, deviceID, key, twin, msgTwin, dealType)
diff --git a/edge/pkg/devicetwin/dtmanager/twin_test.go b/edge/pkg/devicetwin/dtmanager/twin_test.go
index f7d5b43d3..a8eb94dca 100644
--- a/edge/pkg/devicetwin/dtmanager/twin_test.go
+++ b/edge/pkg/devicetwin/dtmanager/twin_test.go
@@ -786,7 +786,7 @@ func TestDealTwinDelete(t *testing.T) {
twin *dttype.MsgTwin
msgTwin *dttype.MsgTwin
dealType int
- err error
+ want *dttype.DealTwinResult
}{
{
name: "TestDealTwinDelete(): Case 1: msgTwin is not nil; isChange is false",
@@ -819,7 +819,11 @@ func TestDealTwinDelete(t *testing.T) {
ActualVersion: &dttype.TwinVersion{},
},
dealType: SyncDealType,
- err: nil,
+ want: &dttype.DealTwinResult{
+ Document: doc,
+ SyncResult: sync,
+ Result: result,
+ },
},
{
name: "TestDealTwinDelete(): Case 2: hasTwinExpected is true; dealVersion() returns false",
@@ -856,7 +860,21 @@ func TestDealTwinDelete(t *testing.T) {
ActualVersion: &dttype.TwinVersion{},
},
dealType: SyncDealType,
- err: nil,
+ want: &dttype.DealTwinResult{
+ Document: make(map[string]*dttype.TwinDoc),
+ SyncResult: map[string]*dttype.MsgTwin{
+ key1: {
+ Optional: &optionTrue,
+ Metadata: &dttype.TypeMetadata{
+ Type: typeString,
+ },
+ ExpectedVersion: &dttype.TwinVersion{
+ CloudVersion: 1,
+ },
+ },
+ },
+ Result: result,
+ },
},
{
name: "TestDealTwinDelete(): Case 3: hasTwinActual is true; dealVersion() returns false",
@@ -893,7 +911,22 @@ func TestDealTwinDelete(t *testing.T) {
},
},
dealType: SyncDealType,
- err: nil,
+ want: &dttype.DealTwinResult{
+ Document: make(map[string]*dttype.TwinDoc),
+ SyncResult: map[string]*dttype.MsgTwin{
+ key1: {
+ Optional: &optionTrue,
+ Metadata: &dttype.TypeMetadata{
+ Type: typeString,
+ },
+ ActualVersion: &dttype.TwinVersion{
+ CloudVersion: 1,
+ },
+ ExpectedVersion: &dttype.TwinVersion{},
+ },
+ },
+ Result: result,
+ },
},
{
name: "TestDealTwinDelete(): Case 4: hasTwinExpected is true; hasTwinActual is true",
@@ -909,7 +942,49 @@ func TestDealTwinDelete(t *testing.T) {
ActualVersion: &dttype.TwinVersion{},
},
dealType: RestDealType,
- err: nil,
+ want: &dttype.DealTwinResult{
+ Document: doc,
+ SyncResult: map[string]*dttype.MsgTwin{
+ key1: {
+ Optional: &optionTrue,
+ Metadata: &dttype.TypeMetadata{
+ Type: dtcommon.TypeDeleted,
+ },
+ Expected: &dttype.TwinValue{
+ Value: nil,
+ Metadata: nil,
+ },
+ ExpectedVersion: &dttype.TwinVersion{
+ CloudVersion: 0,
+ EdgeVersion: 1,
+ },
+ ActualVersion: &dttype.TwinVersion{
+ CloudVersion: 0,
+ EdgeVersion: 1,
+ },
+ Actual: &dttype.TwinValue{
+ Value: nil,
+ Metadata: nil,
+ },
+ },
+ },
+ Result: map[string]*dttype.MsgTwin{
+ key1: nil,
+ },
+ Update: []dtclient.DeviceTwinUpdate{{
+ DeviceID: deviceA,
+ Name: key1,
+ Cols: map[string]interface{}{
+ "attr_type": dtcommon.TypeDeleted,
+ "expected_meta": nil,
+ "expected": nil,
+ "actual_meta": nil,
+ "actual": nil,
+ "actual_version": "{\"cloud\": 0,\"edge\": 1}",
+ "expected_version": "{\"cloud\": 0,\"edge\": 1}",
+ },
+ }},
+ },
},
{
name: "TestDealTwinDelete(): Case 5: hasTwinExpected is true; hasTwinActual is false",
@@ -942,13 +1017,32 @@ func TestDealTwinDelete(t *testing.T) {
ActualVersion: &dttype.TwinVersion{},
},
dealType: SyncDealType,
- err: nil,
+ want: &dttype.DealTwinResult{
+ Document: doc,
+ SyncResult: map[string]*dttype.MsgTwin{
+ key1: nil,
+ },
+ Update: []dtclient.DeviceTwinUpdate{{
+ DeviceID: deviceA,
+ Name: key1,
+ Cols: map[string]interface{}{
+ "attr_type": dtcommon.TypeDeleted,
+ "expected_meta": nil,
+ "expected": nil,
+ "expected_version": "{\"cloud\": 0,\"edge\": 0}",
+ },
+ }},
+ Result: map[string]*dttype.MsgTwin{
+ key1: nil,
+ },
+ },
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
- if err := dealTwinDelete(test.returnResult, test.deviceID, test.key, test.twin, test.msgTwin, test.dealType); !reflect.DeepEqual(err, test.err) {
- t.Errorf("DTManager.TestDealTwinDelete() case failed: got = %+v, Want = %+v", err, test.err)
+ dealTwinDelete(test.returnResult, test.deviceID, test.key, test.twin, test.msgTwin, test.dealType)
+ if (test.returnResult.SyncResult[key1] != nil || test.want.SyncResult[key1] != nil) && !reflect.DeepEqual(*test.returnResult.SyncResult[key1].ExpectedVersion, *test.want.SyncResult[key1].ExpectedVersion) {
+ t.Errorf("DTManager.TestDealTwinDelete() case failed: SyncResult got = %+v, Want = %+v", *test.returnResult.SyncResult[key1].ExpectedVersion, *test.want.SyncResult[key1].ExpectedVersion)
}
})
}
diff --git a/edge/pkg/edged/edged_pods.go b/edge/pkg/edged/edged_pods.go
index 3912cd332..91191baf4 100644
--- a/edge/pkg/edged/edged_pods.go
+++ b/edge/pkg/edged/edged_pods.go
@@ -638,7 +638,7 @@ func (e *edged) GenerateRunContainerOptions(pod *v1.Pod, container *v1.Container
volumes := e.volumeManager.GetMountedVolumesForPod(podName)
blkutil := volumepathhandler.NewBlockVolumePathHandler()
- blkVolumes, err := e.makeBlockVolumes(pod, container, volumes, blkutil)
+ blkVolumes, err := e.makeBlockVolumes(container, volumes, blkutil)
if err != nil {
return nil, nil, err
}
@@ -884,7 +884,7 @@ func (e *edged) podFieldSelectorRuntimeValue(fs *v1.ObjectFieldSelector, pod *v1
// makeBlockVolumes maps the raw block devices specified in the path of the container
// Experimental
-func (e *edged) makeBlockVolumes(pod *v1.Pod, container *v1.Container, podVolumes kubecontainer.VolumeMap, blkutil volumepathhandler.BlockVolumePathHandler) ([]kubecontainer.DeviceInfo, error) {
+func (e *edged) makeBlockVolumes(container *v1.Container, podVolumes kubecontainer.VolumeMap, blkutil volumepathhandler.BlockVolumePathHandler) ([]kubecontainer.DeviceInfo, error) {
var devices []kubecontainer.DeviceInfo
for _, device := range container.VolumeDevices {
// check path is absolute
diff --git a/edge/pkg/edged/edged_status.go b/edge/pkg/edged/edged_status.go
index 9624c86e2..9a9faa5e1 100644
--- a/edge/pkg/edged/edged_status.go
+++ b/edge/pkg/edged/edged_status.go
@@ -116,10 +116,7 @@ func (e *edged) initialNode() (*v1.Node, error) {
return nil, err
}
- err = e.setCPUInfo(node.Status.Capacity, node.Status.Allocatable)
- if err != nil {
- return nil, err
- }
+ e.setCPUInfo(node.Status.Capacity, node.Status.Allocatable)
err = e.setEphemeralStorageInfo(node.Status.Capacity, node.Status.Allocatable)
if err != nil {
@@ -356,11 +353,9 @@ func (e *edged) setMemInfo(total, allocatable v1.ResourceList) error {
return nil
}
-func (e *edged) setCPUInfo(total, allocatable v1.ResourceList) error {
+func (e *edged) setCPUInfo(total, allocatable v1.ResourceList) {
total[v1.ResourceCPU] = resource.MustParse(fmt.Sprintf("%d", runtime.NumCPU()))
allocatable[v1.ResourceCPU] = total[v1.ResourceCPU].DeepCopy()
-
- return nil
}
func (e *edged) setEphemeralStorageInfo(total, allocatable v1.ResourceList) error {
diff --git a/edge/pkg/edged/volume/csi/csi_attacher.go b/edge/pkg/edged/volume/csi/csi_attacher.go
index 6524edfb0..efa60f6f9 100644
--- a/edge/pkg/edged/volume/csi/csi_attacher.go
+++ b/edge/pkg/edged/volume/csi/csi_attacher.go
@@ -75,7 +75,7 @@ func (c *csiAttacher) Attach(spec *volume.Spec, nodeName types.NodeName) (string
func (c *csiAttacher) WaitForAttach(spec *volume.Spec, _ string, pod *v1.Pod, timeout time.Duration) (string, error) {
source, err := getPVSourceFromSpec(spec)
if err != nil {
- klog.Error(log("attacher.WaitForAttach failed to extract CSI volume source: %v", err))
+ klog.Error(log("attacher.WaitForAttach failed to extract CSI volume source of pod %v error: %v", pod.Name, err))
return "", err
}
@@ -87,14 +87,13 @@ func (c *csiAttacher) WaitForAttach(spec *volume.Spec, _ string, pod *v1.Pod, ti
func (c *csiAttacher) waitForVolumeAttachment(volumeHandle, attachID string, timeout time.Duration) (string, error) {
klog.V(4).Info(log("probing for updates from CSI driver for [attachment.ID=%v]", attachID))
- err := wait.PollImmediate(time.Second*5, time.Minute*5, func() (bool, error) {
+ err := wait.PollImmediate(time.Second*5, timeout, func() (bool, error) {
klog.V(4).Info(log("probing VolumeAttachment [id=%v]", attachID))
attach, err := c.k8s.StorageV1().VolumeAttachments().Get(context.Background(), attachID, meta.GetOptions{})
if err != nil {
return false, fmt.Errorf("volume %v has GET error for volume attachment %v: %v", volumeHandle, attachID, err)
}
- successful, err := verifyAttachmentStatus(attach, volumeHandle)
- return successful, err
+ return verifyAttachmentStatus(attach, volumeHandle)
})
if err != nil {
return "", err
diff --git a/edge/pkg/edged/volume/csi/csi_plugin.go b/edge/pkg/edged/volume/csi/csi_plugin.go
index 2a1663427..5a2d2842e 100644
--- a/edge/pkg/edged/volume/csi/csi_plugin.go
+++ b/edge/pkg/edged/volume/csi/csi_plugin.go
@@ -38,9 +38,7 @@ import (
api "k8s.io/api/core/v1"
meta "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
- utilruntime "k8s.io/apimachinery/pkg/util/runtime"
utilversion "k8s.io/apimachinery/pkg/util/version"
- "k8s.io/apimachinery/pkg/util/wait"
utilfeature "k8s.io/apiserver/pkg/util/feature"
clientset "k8s.io/client-go/kubernetes"
storagelisters "k8s.io/client-go/listers/storage/v1"
@@ -114,7 +112,7 @@ func (h *RegistrationHandler) ValidatePlugin(pluginName string, endpoint string,
klog.Infof(log("Trying to validate a new CSI Driver with name: %s endpoint: %s versions: %s",
pluginName, endpoint, strings.Join(versions, ",")))
- _, err := h.validateVersions("ValidatePlugin", pluginName, endpoint, versions)
+ _, err := h.validateVersions("ValidatePlugin", pluginName, versions)
if err != nil {
return fmt.Errorf("validation failed for CSI Driver %s at endpoint %s: %v", pluginName, endpoint, err)
}
@@ -126,7 +124,7 @@ func (h *RegistrationHandler) ValidatePlugin(pluginName string, endpoint string,
func (h *RegistrationHandler) RegisterPlugin(pluginName string, endpoint string, versions []string) error {
klog.Infof(log("Register new plugin with name: %s at endpoint: %s", pluginName, endpoint))
- highestSupportedVersion, err := h.validateVersions("RegisterPlugin", pluginName, endpoint, versions)
+ highestSupportedVersion, err := h.validateVersions("RegisterPlugin", pluginName, versions)
if err != nil {
return err
}
@@ -166,7 +164,7 @@ func (h *RegistrationHandler) RegisterPlugin(pluginName string, endpoint string,
return nil
}
-func (h *RegistrationHandler) validateVersions(callerName, pluginName string, endpoint string, versions []string) (*utilversion.Version, error) {
+func (h *RegistrationHandler) validateVersions(callerName, pluginName string, versions []string) (*utilversion.Version, error) {
if len(versions) == 0 {
err := fmt.Errorf("%s for CSI driver %q failed. Plugin returned an empty list for supported versions", callerName, pluginName)
klog.Error(err)
@@ -254,62 +252,11 @@ func (p *csiPlugin) Init(host volume.VolumeHost) error {
return nil
}
-func initializeCSINode(host volume.VolumeHost) error {
- kvh, ok := host.(volume.KubeletVolumeHost)
- if !ok {
- klog.V(4).Info("Cast from VolumeHost to KubeletVolumeHost failed. Skipping CSINodeInfo initialization, not running on kubelet")
- return nil
- }
- kubeClient := host.GetKubeClient()
- if kubeClient == nil {
- // Kubelet running in standalone mode. Skip CSINodeInfo initialization
- klog.Warning("Skipping CSINodeInfo initialization, kubelet running in standalone mode")
- return nil
- }
-
- kvh.SetKubeletError(errors.New("CSINodeInfo is not yet initialized"))
-
- go func() {
- defer utilruntime.HandleCrash()
-
- // Backoff parameters tuned to retry over 140 seconds. Will fail and restart the Kubelet
- // after max retry steps.
- initBackoff := wait.Backoff{
- Steps: 6,
- Duration: 15 * time.Millisecond,
- Factor: 6.0,
- Jitter: 0.1,
- }
- err := wait.ExponentialBackoff(initBackoff, func() (bool, error) {
- klog.V(4).Infof("Initializing migrated drivers on CSINodeInfo")
- err := nim.InitializeCSINodeWithAnnotation()
- if err != nil {
- kvh.SetKubeletError(fmt.Errorf("failed to initialize CSINodeInfo: %v", err))
- klog.Errorf("Failed to initialize CSINodeInfo: %v", err)
- return false, nil
- }
-
- // Successfully initialized drivers, allow Kubelet to post Ready
- kvh.SetKubeletError(nil)
- return true, nil
- })
- if err != nil {
- // 2 releases after CSIMigration and all CSIMigrationX (where X is a volume plugin)
- // are permanently enabled the apiserver/controllers can assume that the kubelet is
- // using CSI for all Migrated volume plugins. Then all the CSINode initialization
- // code can be dropped from Kubelet.
- // Kill the Kubelet process and allow it to restart to retry initialization
- klog.Fatalf("Failed to initialize CSINodeInfo after retrying")
- }
- }()
- return nil
-}
-
func (p *csiPlugin) GetPluginName() string {
return CSIPluginName
}
-// GetvolumeName returns a concatenated string of CSIVolumeSource.Driver<volNameSe>CSIVolumeSource.VolumeHandle
+// GetVolumeName returns a concatenated string of CSIVolumeSource.Driver<volNameSe>CSIVolumeSource.VolumeHandle
// That string value is used in Detach() to extract driver name and volumeName.
func (p *csiPlugin) GetVolumeName(spec *volume.Spec) (string, error) {
csi, err := getPVSourceFromSpec(spec)
diff --git a/edgesite/cmd/edgesite-server/app/server.go b/edgesite/cmd/edgesite-server/app/server.go
index 3b8034cb8..d0e1acc63 100644
--- a/edgesite/cmd/edgesite-server/app/server.go
+++ b/edgesite/cmd/edgesite-server/app/server.go
@@ -105,15 +105,10 @@ func (p *Proxy) run(o *options.ProxyRunOptions) error {
return fmt.Errorf("failed to run the agent server: %v", err)
}
klog.V(1).Infoln("Starting admin server for debug connections.")
- err = p.runAdminServer(o, server)
- if err != nil {
- return fmt.Errorf("failed to run the admin server: %v", err)
- }
+ p.runAdminServer(o)
+
klog.V(1).Infoln("Starting health server for healthchecks.")
- err = p.runHealthServer(o, server)
- if err != nil {
- return fmt.Errorf("failed to run the health server: %v", err)
- }
+ p.runHealthServer(o, server)
stopCh := SetupSignalHandler()
<-stopCh
@@ -326,7 +321,7 @@ func (p *Proxy) runAgentServer(o *options.ProxyRunOptions, server *server.ProxyS
return nil
}
-func (p *Proxy) runAdminServer(o *options.ProxyRunOptions, server *server.ProxyServer) error {
+func (p *Proxy) runAdminServer(o *options.ProxyRunOptions) {
muxHandler := http.NewServeMux()
muxHandler.Handle("/metrics", promhttp.Handler())
if o.EnableProfiling {
@@ -349,11 +344,9 @@ func (p *Proxy) runAdminServer(o *options.ProxyRunOptions, server *server.ProxyS
}
klog.V(1).Infoln("Admin server stopped listening")
}()
-
- return nil
}
-func (p *Proxy) runHealthServer(o *options.ProxyRunOptions, server *server.ProxyServer) error {
+func (p *Proxy) runHealthServer(o *options.ProxyRunOptions, server *server.ProxyServer) {
livenessHandler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
fmt.Fprintf(w, "ok")
})
@@ -384,6 +377,4 @@ func (p *Proxy) runHealthServer(o *options.ProxyRunOptions, server *server.Proxy
}
klog.V(1).Infoln("Health server stopped listening")
}()
-
- return nil
}
diff --git a/keadm/cmd/keadm/app/cmd/util/common.go b/keadm/cmd/keadm/app/cmd/util/common.go
index a652b4ce0..c97c5051b 100644
--- a/keadm/cmd/keadm/app/cmd/util/common.go
+++ b/keadm/cmd/keadm/app/cmd/util/common.go
@@ -451,7 +451,7 @@ func installKubeEdge(options types.InstallOptions, version semver.Version) error
}
// runEdgeCore starts edgecore with logs being captured
-func runEdgeCore(version semver.Version) error {
+func runEdgeCore() error {
// create the log dir for kubeedge
err := os.MkdirAll(KubeEdgeLogPath, os.ModePerm)
if err != nil {
diff --git a/keadm/cmd/keadm/app/cmd/util/debinstaller.go b/keadm/cmd/keadm/app/cmd/util/debinstaller.go
index 62b791990..270065b31 100644
--- a/keadm/cmd/keadm/app/cmd/util/debinstaller.go
+++ b/keadm/cmd/keadm/app/cmd/util/debinstaller.go
@@ -77,7 +77,7 @@ func (d *DebOS) InstallKubeEdge(options types.InstallOptions) error {
// RunEdgeCore starts edgecore with logs being captured
func (d *DebOS) RunEdgeCore() error {
- return runEdgeCore(d.KubeEdgeVersion)
+ return runEdgeCore()
}
// KillKubeEdgeBinary will search for KubeEdge process and forcefully kill it
diff --git a/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go b/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go
index 9def84768..d4cc8a346 100644
--- a/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go
+++ b/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go
@@ -60,7 +60,7 @@ func (o *PacmanOS) InstallKubeEdge(options types.InstallOptions) error {
// RunEdgeCore sets the environment variable GOARCHAIUS_CONFIG_PATH for the configuration path
// and the starts edgecore with logs being captured
func (o *PacmanOS) RunEdgeCore() error {
- return runEdgeCore(o.KubeEdgeVersion)
+ return runEdgeCore()
}
// KillKubeEdgeBinary will search for KubeEdge process and forcefully kill it
diff --git a/keadm/cmd/keadm/app/cmd/util/rpminstaller.go b/keadm/cmd/keadm/app/cmd/util/rpminstaller.go
index b07abb036..edf67588e 100644
--- a/keadm/cmd/keadm/app/cmd/util/rpminstaller.go
+++ b/keadm/cmd/keadm/app/cmd/util/rpminstaller.go
@@ -98,7 +98,7 @@ func (r *RpmOS) InstallKubeEdge(options types.InstallOptions) error {
// RunEdgeCore starts edgecore with logs being captured
func (r *RpmOS) RunEdgeCore() error {
- return runEdgeCore(r.KubeEdgeVersion)
+ return runEdgeCore()
}
// KillKubeEdgeBinary will search for KubeEdge process and forcefully kill it
diff --git a/staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go b/staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go
index 9ca103bd5..067336f76 100644
--- a/staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go
+++ b/staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go
@@ -26,7 +26,7 @@ var (
autoFree = false
)
-// the connection based on quic protocol
+// QuicConnection the connection based on quic protocol
type QuicConnection struct {
writeDeadline time.Time
readDeadline time.Time
@@ -42,7 +42,7 @@ type QuicConnection struct {
autoRoute bool
}
-// new quic connection
+// NewQuicConn new quic connection
func NewQuicConn(options *ConnectionOptions) *QuicConnection {
quicSession := options.Base.(quic.Session)
return &QuicConnection{
@@ -71,16 +71,6 @@ func (conn *QuicConnection) headerMessage(msg *model.Message) error {
return nil
}
-// process control messages
-func (conn *QuicConnection) processControlMessage(msg *model.Message) error {
- switch msg.GetOperation() {
- case comm.ControlTypeConfig:
- case comm.ControlTypePing:
- case comm.ControlTypePong:
- }
- return nil
-}
-
// read control message from control lan and process
func (conn *QuicConnection) serveControlLan() {
var msg model.Message
@@ -92,15 +82,8 @@ func (conn *QuicConnection) serveControlLan() {
return
}
- // process control message
- result := comm.RespTypeAck
- err = conn.processControlMessage(&msg)
- if err != nil {
- result = comm.RespTypeNack
- }
-
// feedback the response
- resp := msg.NewRespByMessage(&msg, result)
+ resp := msg.NewRespByMessage(&msg, comm.RespTypeAck)
err = conn.ctrlLan.WriteMessage(resp)
if err != nil {
klog.Errorf("failed to send response back, error:%+v", err)
diff --git a/staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go b/staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go
index ed6a586ad..827823568 100644
--- a/staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go
+++ b/staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go
@@ -56,16 +56,6 @@ func (conn *WSConnection) ServeConn() {
}
}
-// process control messages
-func (conn *WSConnection) processControlMessage(msg *model.Message) error {
- switch msg.GetOperation() {
- case comm.ControlTypeConfig:
- case comm.ControlTypePing:
- case comm.ControlTypePong:
- }
- return nil
-}
-
func (conn *WSConnection) filterControlMessage(msg *model.Message) bool {
// check control message
operation := msg.GetOperation()
@@ -75,17 +65,10 @@ func (conn *WSConnection) filterControlMessage(msg *model.Message) bool {
return false
}
- // process control message
- result := comm.RespTypeAck
- err := conn.processControlMessage(msg)
- if err != nil {
- result = comm.RespTypeNack
- }
-
// feedback the response
- resp := msg.NewRespByMessage(msg, result)
+ resp := msg.NewRespByMessage(msg, comm.RespTypeAck)
conn.locker.Lock()
- err = lane.NewLane(api.ProtocolTypeWS, conn.wsConn).WriteMessage(resp)
+ err := lane.NewLane(api.ProtocolTypeWS, conn.wsConn).WriteMessage(resp)
conn.locker.Unlock()
if err != nil {
klog.Errorf("failed to send response back, error:%+v", err)