diff options
| author | cl2017 <chenlin.liu@daocloud.io> | 2023-09-21 15:08:46 +0800 |
|---|---|---|
| committer | cl2017 <chenlin.liu@daocloud.io> | 2023-09-22 13:00:03 +0800 |
| commit | b608a115782ef895ec8fec5212982cad1da157d1 (patch) | |
| tree | aaad6c9a0891ae4ca8b30a4aaf957ca172bc80d5 /edge | |
| parent | device crd v1beta1 cloud (diff) | |
| download | kubeedge-b608a115782ef895ec8fec5212982cad1da157d1.tar.gz | |
device crd v1beta1 edge
Signed-off-by: cl2017 <chenlin.liu@daocloud.io>
Diffstat (limited to 'edge')
| -rw-r--r-- | edge/pkg/devicetwin/dmiclient/client.go | 39 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dmiserver/server.go | 30 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtcommon/util.go | 26 | ||||
| -rw-r--r-- | edge/pkg/devicetwin/dtmanager/dmiworker.go | 16 |
4 files changed, 37 insertions, 74 deletions
diff --git a/edge/pkg/devicetwin/dmiclient/client.go b/edge/pkg/devicetwin/dmiclient/client.go index 555678a46..af6e78a10 100644 --- a/edge/pkg/devicetwin/dmiclient/client.go +++ b/edge/pkg/devicetwin/dmiclient/client.go @@ -28,8 +28,8 @@ import ( deviceconst "github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller/constants" "github.com/kubeedge/kubeedge/edge/pkg/devicetwin/dtcommon" - "github.com/kubeedge/kubeedge/pkg/apis/devices/v1alpha2" - dmiapi "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/apis/devices/v1beta1" + dmiapi "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1beta1" ) type DMIClient struct { @@ -82,7 +82,7 @@ func (dc *DMIClient) close() { dc.CancelFunc() } -func createDeviceRequest(device *v1alpha2.Device) (*dmiapi.RegisterDeviceRequest, error) { +func createDeviceRequest(device *v1beta1.Device) (*dmiapi.RegisterDeviceRequest, error) { d, err := dtcommon.ConvertDevice(device) if err != nil { return nil, err @@ -99,7 +99,7 @@ func removeDeviceRequest(deviceName string) (*dmiapi.RemoveDeviceRequest, error) }, nil } -func updateDeviceRequest(device *v1alpha2.Device) (*dmiapi.UpdateDeviceRequest, error) { +func updateDeviceRequest(device *v1beta1.Device) (*dmiapi.UpdateDeviceRequest, error) { d, err := dtcommon.ConvertDevice(device) if err != nil { return nil, err @@ -110,7 +110,7 @@ func updateDeviceRequest(device *v1alpha2.Device) (*dmiapi.UpdateDeviceRequest, }, nil } -func createDeviceModelRequest(model *v1alpha2.DeviceModel) (*dmiapi.CreateDeviceModelRequest, error) { +func createDeviceModelRequest(model *v1beta1.DeviceModel) (*dmiapi.CreateDeviceModelRequest, error) { m, err := dtcommon.ConvertDeviceModel(model) if err != nil { return nil, err @@ -121,7 +121,7 @@ func createDeviceModelRequest(model *v1alpha2.DeviceModel) (*dmiapi.CreateDevice }, nil } -func updateDeviceModelRequest(model *v1alpha2.DeviceModel) (*dmiapi.UpdateDeviceModelRequest, error) { +func updateDeviceModelRequest(model *v1beta1.DeviceModel) (*dmiapi.UpdateDeviceModelRequest, error) { m, err := dtcommon.ConvertDeviceModel(model) if err != nil { return nil, err @@ -179,11 +179,8 @@ func (dcs *DMIClients) getDMIClientConn(protocol string) (*DMIClient, error) { return dc, nil } -func (dcs *DMIClients) RegisterDevice(device *v1alpha2.Device) error { - protocol, err := dtcommon.GetProtocolNameOfDevice(device) - if err != nil { - return err - } +func (dcs *DMIClients) RegisterDevice(device *v1beta1.Device) error { + protocol := device.Spec.Protocol.ProtocolName dc, err := dcs.getDMIClientConn(protocol) if err != nil { @@ -203,11 +200,8 @@ func (dcs *DMIClients) RegisterDevice(device *v1alpha2.Device) error { return nil } -func (dcs *DMIClients) RemoveDevice(device *v1alpha2.Device) error { - protocol, err := dtcommon.GetProtocolNameOfDevice(device) - if err != nil { - return err - } +func (dcs *DMIClients) RemoveDevice(device *v1beta1.Device) error { + protocol := device.Spec.Protocol.ProtocolName dc, err := dcs.getDMIClientConn(protocol) if err != nil { @@ -227,11 +221,8 @@ func (dcs *DMIClients) RemoveDevice(device *v1alpha2.Device) error { return nil } -func (dcs *DMIClients) UpdateDevice(device *v1alpha2.Device) error { - protocol, err := dtcommon.GetProtocolNameOfDevice(device) - if err != nil { - return err - } +func (dcs *DMIClients) UpdateDevice(device *v1beta1.Device) error { + protocol := device.Spec.Protocol.ProtocolName dc, err := dcs.getDMIClientConn(protocol) if err != nil { @@ -251,7 +242,7 @@ func (dcs *DMIClients) UpdateDevice(device *v1alpha2.Device) error { return nil } -func (dcs *DMIClients) CreateDeviceModel(model *v1alpha2.DeviceModel) error { +func (dcs *DMIClients) CreateDeviceModel(model *v1beta1.DeviceModel) error { protocol := model.Spec.Protocol dc, err := dcs.getDMIClientConn(protocol) if err != nil { @@ -271,7 +262,7 @@ func (dcs *DMIClients) CreateDeviceModel(model *v1alpha2.DeviceModel) error { return nil } -func (dcs *DMIClients) RemoveDeviceModel(model *v1alpha2.DeviceModel) error { +func (dcs *DMIClients) RemoveDeviceModel(model *v1beta1.DeviceModel) error { protocol := model.Spec.Protocol dc, err := dcs.getDMIClientConn(protocol) if err != nil { @@ -291,7 +282,7 @@ func (dcs *DMIClients) RemoveDeviceModel(model *v1alpha2.DeviceModel) error { return nil } -func (dcs *DMIClients) UpdateDeviceModel(model *v1alpha2.DeviceModel) error { +func (dcs *DMIClients) UpdateDeviceModel(model *v1beta1.DeviceModel) error { protocol := model.Spec.Protocol dc, err := dcs.getDMIClientConn(protocol) if err != nil { diff --git a/edge/pkg/devicetwin/dmiserver/server.go b/edge/pkg/devicetwin/dmiserver/server.go index a63d20ae8..cd5516e82 100644 --- a/edge/pkg/devicetwin/dmiserver/server.go +++ b/edge/pkg/devicetwin/dmiserver/server.go @@ -41,8 +41,8 @@ import ( "github.com/kubeedge/kubeedge/edge/pkg/devicetwin/dmiclient" "github.com/kubeedge/kubeedge/edge/pkg/devicetwin/dtcommon" "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao" - "github.com/kubeedge/kubeedge/pkg/apis/devices/v1alpha2" - pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/apis/devices/v1beta1" + pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1beta1" ) const ( @@ -61,8 +61,8 @@ type DMICache struct { DeviceMu *sync.Mutex DeviceModelMu *sync.Mutex MapperList map[string]*pb.MapperInfo - DeviceModelList map[string]*v1alpha2.DeviceModel - DeviceList map[string]*v1alpha2.Device + DeviceModelList map[string]*v1beta1.DeviceModel + DeviceList map[string]*v1beta1.Device } func (s *server) MapperRegister(ctx context.Context, in *pb.MapperRegisterRequest) (*pb.MapperRegisterResponse, error) { @@ -94,11 +94,7 @@ func (s *server) MapperRegister(ctx context.Context, in *pb.MapperRegisterReques s.dmiCache.DeviceMu.Lock() defer s.dmiCache.DeviceMu.Unlock() for _, device := range s.dmiCache.DeviceList { - protocol, err := dtcommon.GetProtocolNameOfDevice(device) - if err != nil { - klog.Errorf("fail to get protocol name with err: %+v", err) - continue - } + protocol := device.Spec.Protocol.ProtocolName if protocol == in.Mapper.Protocol { dev, err := dtcommon.ConvertDevice(device) @@ -137,13 +133,7 @@ func (s *server) ReportDeviceStatus(ctx context.Context, in *pb.ReportDeviceStat } for _, twin := range in.ReportedDevice.Twins { - propertyType, ok := twin.Reported.Metadata[PropertyType] - if !ok { - errLog := fmt.Sprintf("fail to get propertyType for property %s of device %s", twin.PropertyName, in.DeviceName) - klog.Errorf(errLog) - return nil, fmt.Errorf(errLog) - } - msg, err := CreateMessageTwinUpdate(twin.PropertyName, propertyType, twin.Reported.Value) + msg, err := CreateMessageTwinUpdate(twin) if err != nil { klog.Errorf("fail to create message data for property %s of device %s with err: %v", twin.PropertyName, in.DeviceName, err) return nil, err @@ -166,14 +156,14 @@ func handleDeviceTwin(deviceName string, payload []byte) { } // CreateMessageTwinUpdate create twin update message. -func CreateMessageTwinUpdate(name, valueType, value string) ([]byte, error) { +func CreateMessageTwinUpdate(twin *pb.Twin) ([]byte, error) { var updateMsg DeviceTwinUpdate updateMsg.BaseMessage.Timestamp = getTimestamp() updateMsg.Twin = map[string]*types.MsgTwin{} - updateMsg.Twin[name] = &types.MsgTwin{} - updateMsg.Twin[name].Actual = &types.TwinValue{Value: &value} - updateMsg.Twin[name].Metadata = &types.TypeMetadata{Type: valueType} + updateMsg.Twin[twin.PropertyName] = &types.MsgTwin{} + updateMsg.Twin[twin.PropertyName].Expected = &types.TwinValue{Value: &twin.ObservedDesired.Value} + updateMsg.Twin[twin.PropertyName].Actual = &types.TwinValue{Value: &twin.Reported.Value} msg, err := json.Marshal(updateMsg) return msg, err diff --git a/edge/pkg/devicetwin/dtcommon/util.go b/edge/pkg/devicetwin/dtcommon/util.go index 2adb4d373..eff5156aa 100644 --- a/edge/pkg/devicetwin/dtcommon/util.go +++ b/edge/pkg/devicetwin/dtcommon/util.go @@ -3,7 +3,6 @@ package dtcommon import ( "encoding/json" "errors" - "fmt" "regexp" "strconv" "strings" @@ -11,8 +10,8 @@ import ( "k8s.io/klog/v2" "github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller/constants" - "github.com/kubeedge/kubeedge/pkg/apis/devices/v1alpha2" - pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/apis/devices/v1beta1" + pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1beta1" ) // ValidateValue validate value type @@ -61,23 +60,7 @@ func ValidateTwinValue(value string) bool { return match } -func GetProtocolNameOfDevice(device *v1alpha2.Device) (string, error) { - protocol := device.Spec.Protocol - if protocol.OpcUA != nil { - return constants.OPCUA, nil - } - if protocol.Modbus != nil { - return constants.Modbus, nil - } - if protocol.Bluetooth != nil { - return constants.Bluetooth, nil - } - if protocol.CustomizedProtocol != nil { - return protocol.CustomizedProtocol.ProtocolName, nil - } - return "", fmt.Errorf("cannot find protocol name for device %s", device.Name) -} -func ConvertDevice(device *v1alpha2.Device) (*pb.Device, error) { +func ConvertDevice(device *v1beta1.Device) (*pb.Device, error) { data, err := json.Marshal(device) if err != nil { klog.Errorf("fail to marshal device %s with err: %v", device.Name, err) @@ -90,14 +73,13 @@ func ConvertDevice(device *v1alpha2.Device) (*pb.Device, error) { klog.Errorf("fail to unmarshal device %s with err: %v", device.Name, err) return nil, err } - edgeDevice.Name = device.Name edgeDevice.Spec.DeviceModelReference = device.Spec.DeviceModelRef.Name return &edgeDevice, nil } -func ConvertDeviceModel(model *v1alpha2.DeviceModel) (*pb.DeviceModel, error) { +func ConvertDeviceModel(model *v1beta1.DeviceModel) (*pb.DeviceModel, error) { data, err := json.Marshal(model) if err != nil { klog.Errorf("fail to marshal device model %s with err: %v", model.Name, err) diff --git a/edge/pkg/devicetwin/dtmanager/dmiworker.go b/edge/pkg/devicetwin/dtmanager/dmiworker.go index 435e9e179..e46e2513f 100644 --- a/edge/pkg/devicetwin/dtmanager/dmiworker.go +++ b/edge/pkg/devicetwin/dtmanager/dmiworker.go @@ -33,8 +33,8 @@ import ( "github.com/kubeedge/kubeedge/edge/pkg/devicetwin/dtcontext" "github.com/kubeedge/kubeedge/edge/pkg/devicetwin/dttype" "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao" - "github.com/kubeedge/kubeedge/pkg/apis/devices/v1alpha2" - pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/apis/devices/v1beta1" + pb "github.com/kubeedge/kubeedge/pkg/apis/dmi/v1beta1" ) // TwinWorker deal twin event @@ -52,8 +52,8 @@ func (dw *DMIWorker) init() { DeviceMu: &sync.Mutex{}, DeviceModelMu: &sync.Mutex{}, MapperList: make(map[string]*pb.MapperInfo), - DeviceList: make(map[string]*v1alpha2.Device), - DeviceModelList: make(map[string]*v1alpha2.DeviceModel), + DeviceList: make(map[string]*v1beta1.Device), + DeviceModelList: make(map[string]*v1beta1.DeviceModel), } dw.initDMIActionCallBack() @@ -112,8 +112,8 @@ func (dw *DMIWorker) dealMetaDeviceOperation(context *dtcontext.DTContext, resou if len(resources) != 3 { return fmt.Errorf("wrong resources %s", message.Router.Resource) } - var device v1alpha2.Device - var dm v1alpha2.DeviceModel + var device v1beta1.Device + var dm v1beta1.DeviceModel switch resources[1] { case constants.ResourceTypeDevice: err := json.Unmarshal(message.Content.([]byte), &device) @@ -203,7 +203,7 @@ func (dw *DMIWorker) initDeviceModelInfoFromDB() { } for _, meta := range *metas { - deviceModel := v1alpha2.DeviceModel{} + deviceModel := v1beta1.DeviceModel{} if err := json.Unmarshal([]byte(meta), &deviceModel); err != nil { klog.Errorf("fail to unmarshal device model info from db with err: %v", err) return @@ -223,7 +223,7 @@ func (dw *DMIWorker) initDeviceInfoFromDB() { } for _, meta := range *metas { - device := v1alpha2.Device{} + device := v1beta1.Device{} if err := json.Unmarshal([]byte(meta), &device); err != nil { klog.Errorf("fail to unmarshal device info from db with err: %v", err) return |
