summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authorcl2017 <chenlin.liu@daocloud.io>2023-09-21 15:08:46 +0800
committercl2017 <chenlin.liu@daocloud.io>2023-09-22 13:00:03 +0800
commitb608a115782ef895ec8fec5212982cad1da157d1 (patch)
treeaaad6c9a0891ae4ca8b30a4aaf957ca172bc80d5 /edge
parentdevice crd v1beta1 cloud (diff)
downloadkubeedge-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.go39
-rw-r--r--edge/pkg/devicetwin/dmiserver/server.go30
-rw-r--r--edge/pkg/devicetwin/dtcommon/util.go26
-rw-r--r--edge/pkg/devicetwin/dtmanager/dmiworker.go16
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