diff options
| author | vincentgoat <linguohui1@huawei.com> | 2022-06-25 17:21:07 +0800 |
|---|---|---|
| committer | vincentgoat <linguohui1@huawei.com> | 2022-07-12 00:10:34 +0800 |
| commit | 45342cf3898b83d77180ce5de14ffa9530bbb692 (patch) | |
| tree | 4c5e335fb7ebe9857ef3888b43d8c9e6192f5862 | |
| parent | Merge pull request #3980 from gy95/release-1.9 (diff) | |
| download | kubeedge-45342cf3898b83d77180ce5de14ffa9530bbb692.tar.gz | |
fix complex parameters
Signed-off-by: vincentgoat <linguohui1@huawei.com>
| -rw-r--r-- | edge/pkg/devicetwin/dtcommon/common.go | 2 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/communicate.go | 13 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/communicate_test.go | 2 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/membership.go | 10 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/membership_test.go | 13 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/twin.go | 23 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/twin_test.go | 110 | ||||
| -rw-r--r-- | edge/pkg/edged/edged_pods.go | 4 | ||||
| -rw-r--r-- | edge/pkg/edged/edged_status.go | 9 | ||||
| -rw-r--r-- | edge/pkg/edged/volume/csi/csi_attacher.go | 7 | ||||
| -rw-r--r-- | edge/pkg/edged/volume/csi/csi_plugin.go | 61 | ||||
| -rw-r--r-- | edgesite/cmd/edgesite-server/app/server.go | 19 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/util/common.go | 2 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/util/debinstaller.go | 2 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go | 2 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/util/rpminstaller.go | 2 | ||||
| -rw-r--r-- | staging/src/github.com/kubeedge/viaduct/pkg/conn/quic.go | 23 | ||||
| -rw-r--r-- | staging/src/github.com/kubeedge/viaduct/pkg/conn/ws.go | 21 |
18 files changed, 162 insertions, 163 deletions
diff --git a/edge/pkg/devicetwin/dtcommon/common.go b/edge/pkg/devicetwin/dtcommon/common.go index 79498c920..e8aef910f 100644 --- a/edge/pkg/devicetwin/dtcommon/common.go +++ b/edge/pkg/devicetwin/dtcommon/common.go @@ -113,4 +113,6 @@ const ( ConflictCode = 409 // InternalErrorCode server internal error InternalErrorCode = 500 + // TypeDeleted type deleted + TypeDeleted = "deleted" ) 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 baa405bf5..b57ab469c 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 36d6e28ac..c1dd767eb 100644 --- a/edge/pkg/devicetwin/dtmanager/membership_test.go +++ b/edge/pkg/devicetwin/dtmanager/membership_test.go @@ -19,6 +19,7 @@ package dtmanager import ( "encoding/json" "errors" + "reflect" "sync" "testing" @@ -280,7 +281,9 @@ func TestDealMembershipGetValid(t *testing.T) { Content: content, } err := dealMembershipGet(dtc, "t", m) - assert.NoError(t, 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) + } } func TestDealMembershipGetInnerValid(t *testing.T) { @@ -306,7 +309,9 @@ func TestDealMembershipGetInnerValid(t *testing.T) { content, _ := json.Marshal(payload) err := dealMembershipGetInner(dtc, content) - assert.NoError(t, 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) + } } func TestDealMembershipGetInnerInValid(t *testing.T) { @@ -318,7 +323,9 @@ func TestDealMembershipGetInnerInValid(t *testing.T) { } err := dealMembershipGetInner(dtc, []byte("invalid")) - assert.NoError(t, 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) + } } //Commented As we are not considering about the coverage incase for coverage we can uncomment below cases. diff --git a/edge/pkg/devicetwin/dtmanager/twin.go b/edge/pkg/devicetwin/dtmanager/twin.go index 6f2f61e1a..83a57733f 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 @@ -369,7 +369,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,9 +379,9 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st syncResult[key] = &dttype.MsgTwin{} update := returnResult.Update isChange := false - if msgTwin == nil && dealType == RestDealType && *twin.Optional || dealType >= SyncDealType && strings.Compare(msgTwin.Metadata.Type, "deleted") == 0 { - if twin.Metadata != nil && strings.Compare(twin.Metadata.Type, "deleted") == 0 { - return nil + 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 } if dealType != RestDealType { dealType = SyncTwinDeleteDealType @@ -405,7 +405,7 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st syncResult[key] = ©Sync delete(document, key) returnResult.SyncResult = syncResult - return nil + return } } else { expectedVersionJSON, _ := json.Marshal(expectedVersion) @@ -447,7 +447,7 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st syncResult[key] = ©Sync delete(document, key) returnResult.SyncResult = syncResult - return nil + return } } else { actualVersionJSON, _ := json.Marshal(actualVersion) @@ -487,8 +487,6 @@ func dealTwinDelete(returnResult *dttype.DealTwinResult, deviceID string, key st delete(document, key) delete(syncResult, key) } - - return nil } //0:expected ,1 :actual @@ -946,11 +944,8 @@ func DealMsgTwin(context *dtcontext.DTContext, deviceID string, msgTwins map[str if dealType >= 1 && msgTwin != nil && (msgTwin.Metadata == nil) { klog.Infof("Not found metadata of twin") } - if msgTwin == nil && dealType == 0 || dealType >= 1 && strings.Compare(msgTwin.Metadata.Type, "deleted") == 0 { - err = dealTwinDelete(&returnResult, deviceID, key, twin, msgTwin, dealType) - if err != nil { - return returnResult - } + if msgTwin == nil && dealType == 0 || dealType >= 1 && strings.Compare(msgTwin.Metadata.Type, dtcommon.TypeDeleted) == 0 { + 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 34cd23c94..06e3371c1 100644 --- a/edge/pkg/devicetwin/dtmanager/twin_test.go +++ b/edge/pkg/devicetwin/dtmanager/twin_test.go @@ -787,7 +787,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", @@ -820,7 +820,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", @@ -857,7 +861,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", @@ -894,7 +912,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", @@ -910,7 +943,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", @@ -943,13 +1018,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 995a93b01..10d600f7f 100644 --- a/edge/pkg/edged/edged_pods.go +++ b/edge/pkg/edged/edged_pods.go @@ -611,7 +611,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 } @@ -857,7 +857,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 0396fb7c2..e2a4ffad1 100644 --- a/edge/pkg/edged/edged_status.go +++ b/edge/pkg/edged/edged_status.go @@ -117,10 +117,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 { @@ -357,11 +354,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 979e6c30f..2119e5c7e 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 c04196c6f..ca700efa6 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 ed674b1c5..5b7db644d 100644 --- a/edgesite/cmd/edgesite-server/app/server.go +++ b/edgesite/cmd/edgesite-server/app/server.go @@ -106,15 +106,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 @@ -327,7 +322,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 { @@ -350,11 +345,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") }) @@ -385,6 +378,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 fa8525272..6ca48cea2 100644 --- a/keadm/cmd/keadm/app/cmd/util/common.go +++ b/keadm/cmd/keadm/app/cmd/util/common.go @@ -326,7 +326,7 @@ func installKubeEdge(options types.InstallOptions, arch string, version semver.V } // 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 20c2ab43b..4801925d2 100644 --- a/keadm/cmd/keadm/app/cmd/util/debinstaller.go +++ b/keadm/cmd/keadm/app/cmd/util/debinstaller.go @@ -82,7 +82,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 fed0474e1..8acf35e74 100644 --- a/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go +++ b/keadm/cmd/keadm/app/cmd/util/pacmaninstaller.go @@ -77,7 +77,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 0e8fa4b89..f3485ee4e 100644 --- a/keadm/cmd/keadm/app/cmd/util/rpminstaller.go +++ b/keadm/cmd/keadm/app/cmd/util/rpminstaller.go @@ -115,7 +115,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 6746e8594..2a612d052 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 209a75770..30c208866 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) |
