EdgeXFoundry之device-modbus-go代码研究
·
device-modbus-go 代码解析
bin
edgex-launch.sh
目的:
- 启动所有EdgeX Go的二进制文件(必须以前就构建好)
前提:
- Consul和MongoDB都已经安装并运行
清理:
- 杀死进程edgex-device-modbus
流程:
- 先到CMD目录下(…/cmd)
- 执行目录下的device-modbus文件,让其作为edgex-device-modbus进程启动
- 回到PWD目录
- 如果陷入trap,就执行清理功能并退出
- 循环运行监视trap
###
# 启动所有EdgeX Go的二进制文件(必须以前就构建好)
#
# 前提是Consul和MongoDB都已经安装并运行
#
###
DIR=$PWD
CMD=../cmd
function cleanup {
pkill edgex-device-modbus #清除->杀死进程
}
cd $CMD #先到CMD目录下(../cmd)
exec -a edgex-device-modbus ./device-modbus & #执行目录下的device-modbus文件,让其作为edgex-device-modbus进程启动
cd $DIR #回到PWD目录
trap cleanup EXIT #如果陷入trap,就执行清理功能并退出
while : ; do sleep 1 ; done ##循环运行监视trap
test-attribution-txt.sh
环境变量:
- 获取脚本所在目录
- 提取并进入bash脚本第一个参数里的目录并不希望在屏幕上显示输出结果,如果执行成功则将当前目录传给SCRIPT_DIR(脚本目录)
- 提取SCRIPT_DIR的父级目录
清理:
- 进入GIT_ROOT目录,复原到存在vendor的目录下
- 强制删除vendor文件夹
流程:
- 进入GIT_ROOT
- 检查有没有vendor.bk(备份),如果有以1退出,并提示该文件存在,需要先移除再继续
- 如果陷入trap执行清理功能
- 检查有没有vendor.bk,如果有将vendor.bk改为vendor(相当于备份还原)
- 如果vendor文件夹存在,备份他,这样我们就能建一个新的
- go mod时在项目中创建vendor文件将依赖包拷贝过来
- 打开nullglobbing,这样如果cmd目录中没有任何东西,那么我们就不会在这个循环中做任何事情
- 如果找不到Attribution.txt文件,以1退出,并提示该文件找不到
- 循环遍历vendor文件夹中的每个lib文件,确保lib存在且明确,否则提示文件中某lib丢失,请添加
- 回到GIT_ROOT目录
其中Attribution.txt内容是在设备服务modbus-go中引用的开源项目。
#!/bin/bash -e
# 获取脚本所在目录
# 摘自 https://stackoverflow.com/a/246128/10102404
SCRIPT_DIR="$( cd "$( dirname "${BASH_SOURCE[0]}" )" >/dev/null && pwd )" # 提取并进入bash脚本第一个参数里的目录并不希望在屏幕上显示输出结果,如果执行成功则将当前目录传给SCRIPT_DIR(脚本目录)
GIT_ROOT=$(dirname "$SCRIPT_DIR") # 提取SCRIPT_DIR的父级目录
EXIT_CODE=0
cd "$GIT_ROOT" # 进入GIT_ROOT
if [ -d vendor.bk ]; then # 检查有没有vendor.bk(备份),如果有以1退出,并提示该文件存在,需要先移除再继续
echo "vendor.bk exits - remove before continuing"
exit 1
fi
trap cleanup 0 1 2 3 6
cleanup()
{
cd "$GIT_ROOT" # 进入GIT_ROOT目录,复原到存在vendor的目录下
rm -rf vendor # 强制删除vendor文件夹
if [ -d vendor.bk ]; then # 检查有没有vendor.bk,如果有将vendor.bk改为vendor(相当于备份还原)
mv vendor.bk vendor
fi
exit $EXIT_CODE
}
# 如果vendor文件夹存在,备份他,这样我们就能建一个新的
if [ -d vendor ]; then
mv vendor vendor.bk
fi
GO111MODULE=on go mod vendor # 在项目中创建vendor文件将依赖包拷贝过来
# 打开nullglobbing,这样如果cmd目录中没有任何东西,那么我们就不会在这个循环中做任何事情
shopt -s nullglob
if [ ! -f Attribution.txt ]; then # 如果找不到Attribution.txt文件,以1退出,并提示该文件找不到了,请添加
echo "An Attribution.txt file is missing, please add"
EXIT_CODE=1
else
# 循环遍历vendor文件夹中的每个lib文件,确保lib存在且明确,否则提示文件中某lib丢失,请添加
while IFS= read -r lib; do
if ! grep -q "$lib" Attribution.txt && [ "$lib" != "explicit" ]; then
echo "An attribution for $lib is missing from in Attribution.txt, please add"
# need to do this in a bash subshell, see SC2031
(( EXIT_CODE=1 ))
fi
done < <(grep '#' < "$GIT_ROOT/vendor/modules.txt" | awk '{print $2}')
fi
cd "$GIT_ROOT" # 回到GIT_ROOT目录
cmd
main.go
流程:
- 设置服务名称
- 建立对接的协议驱动
- 和SDK对接
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2018-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
package main
import (
"github.com/edgexfoundry/device-sdk-go/v2/pkg/startup"
"github.com/edgexfoundry/device-modbus-go"
"github.com/edgexfoundry/device-modbus-go/internal/driver"
)
const (
serviceName string = "device-modbus" // 设置服务名称
)
func main() {
sd := driver.NewProtocolDriver() // 建立对接的协议驱动
startup.Bootstrap(serviceName, device_modbus.Version, sd) // 和SDK对接
}
internal/driver
config.go
结构体:
- 一个综合了Modbus RTU和Modbus TCP参数的结构体
功能:
- 创建连接信息
- 传入配置协议集
- 检查modbus协议形式
- 对错误配置进行错误提示
- 协议是RTU则调用功能创建RTU类型的连接信息
- 协议是TCP则调用功能创建TCP类型的连接信息
- 解析整数配置参数
- 获取配置参数
- 转换成int类型
- 创建RTU类型连接信息
- 解析RTU协议对应参数初始化ConnectionInfo结构体
- 创建TCP类型连接信息
- 解析TCP协议对应参数初始化ConnectionInfo结构体
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2019-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
package driver
import (
"fmt"
"strconv"
"github.com/edgexfoundry/go-mod-core-contracts/v2/models"
)
// ConnectionInfo 是设备连接所需的信息
type ConnectionInfo struct {// 一个综合了RTU和TCP参数的结构体
Protocol string //协议
Address string //地址
Port int //端口
BaudRate int //波特率
DataBits int //数据位
StopBits int //停止位
Parity string //校验
UnitID uint8 //从站地址
// 连接 & 读取超时(秒)
Timeout int
// 闲置超时(秒) 超时断开连接
IdleTimeout int
}
func createConnectionInfo(protocols map[string]models.ProtocolProperties) (info *ConnectionInfo, err error) { // 传入配置协议集,解析生成对应连接信息
protocolRTU, rtuExist := protocols[ProtocolRTU]// 检查modbus协议形式
protocolTCP, tcpExist := protocols[ProtocolTCP]
if rtuExist && tcpExist {// 对错误配置进行错误提示
return info, fmt.Errorf("unsupported multiple protocols, please choose %s or %s, not both", ProtocolRTU, ProtocolTCP)
} else if !rtuExist && !tcpExist {
return info, fmt.Errorf("unable to create connection info, protocol config '%s' or %s not exist", ProtocolRTU, ProtocolTCP)
}
if rtuExist {// 配置是RTU则创建RTU类型的连接信息
info, err = createRTUConnectionInfo(protocolRTU)
if err != nil {
return nil, err
}
} else if tcpExist {// 配置是TCP则创建TCP类型的连接信息
info, err = createTcpConnectionInfo(protocolTCP)
if err != nil {
return nil, err
}
}
return info, nil
}
func parseIntValue(properties map[string]string, key string) (int, error) { // 将配置解析成int
str, ok := properties[key] // 获取配置参数
if !ok {
return 0, fmt.Errorf("protocol config '%s' not exist", key)
}
val, err := strconv.Atoi(str) // 转换成int类型
if err != nil {
return 0, fmt.Errorf("fail to parse protocol config '%s', %v", key, err)
}
return val, nil
}
func createRTUConnectionInfo(rtuProtocol map[string]string) (info *ConnectionInfo, err error) { // 解析协议参数创建ConnectionInfo结构体
errorMessage := "unable to create RTU connection info, protocol config '%s' not exist"
address, ok := rtuProtocol[Address]
if !ok {
return nil, fmt.Errorf(errorMessage, Address)
}
us, ok := rtuProtocol[UnitID]
if !ok {
return nil, fmt.Errorf(errorMessage, UnitID)
}
unitID, err := strconv.ParseUint(us, 0, 8)
if err != nil {
return nil, fmt.Errorf("uintID value out of range(0–255). Error: %v", err)
}
baudRate, err := parseIntValue(rtuProtocol, BaudRate)
if err != nil {
return nil, err
}
dataBits, err := parseIntValue(rtuProtocol, DataBits)
if err != nil {
return nil, err
}
stopBits, err := parseIntValue(rtuProtocol, StopBits)
if err != nil {
return nil, err
}
parity, ok := rtuProtocol[Parity]
if !ok {
return nil, fmt.Errorf(errorMessage, Parity)
}
if parity != "N" && parity != "O" && parity != "E" {
return nil, fmt.Errorf("invalid parity value, it should be N(None) or O(Odd) or E(Even)")
}
timeout, err := parseIntValue(rtuProtocol, Timeout)
if err != nil {
return nil, err
}
idleTimeout, err := parseIntValue(rtuProtocol, IdleTimeout)
if err != nil {
return nil, err
}
return &ConnectionInfo{
Protocol: ProtocolRTU,
Address: address,
BaudRate: baudRate,
DataBits: dataBits,
StopBits: stopBits,
Parity: parity,
UnitID: byte(unitID),
Timeout: timeout,
IdleTimeout: idleTimeout,
}, nil
}
func createTcpConnectionInfo(tcpProtocol map[string]string) (info *ConnectionInfo, err error) {// 解析协议参数创建ConnectionInfo结构体
errorMessage := "unable to create TCP connection info, protocol config '%s' not exist"
address, ok := tcpProtocol[Address]
if !ok {
return nil, fmt.Errorf(errorMessage, Address)
}
portString, ok := tcpProtocol[Port]
if !ok {
return nil, fmt.Errorf(errorMessage, Port)
}
port, err := strconv.ParseUint(portString, 0, 16)
if err != nil {
return nil, fmt.Errorf("port value out of range(0–65535). Error: %v", err)
}
unitIDString, ok := tcpProtocol[UnitID]
if !ok {
return nil, fmt.Errorf(errorMessage, UnitID)
}
unitID, err := strconv.ParseUint(unitIDString, 0, 8)
if err != nil {
return nil, fmt.Errorf("uintID value out of range(0–255). Error: %v", err)
}
timeout, err := parseIntValue(tcpProtocol, Timeout)
if err != nil {
return nil, err
}
idleTimeout, err := parseIntValue(tcpProtocol, IdleTimeout)
if err != nil {
return nil, err
}
return &ConnectionInfo{
Protocol: ProtocolTCP,
Address: address,
Port: int(port),
UnitID: byte(unitID),
Timeout: timeout,
IdleTimeout: idleTimeout,
}, nil
}
constant.go
常量:
- 数据类型字符化
- 参数名称字符化
map(位数对应表):
- PrimaryTableBitCountMap
- ValueTypeBitCountMap
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2018-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
package driver
import (
"github.com/edgexfoundry/go-mod-core-contracts/v2/common"
)
const (// 数据类型字符化、参数名称字符化
BOOL = "BOOL"
INT16 = "INT16"
INT32 = "INT32"
INT64 = "INT64"
UINT16 = "UINT16"
UINT32 = "UINT32"
UINT64 = "UINT64"
FLOAT32 = "FLOAT32"
FLOAT64 = "FLOAT64"
DISCRETES_INPUT = "DISCRETES_INPUT"
COILS = "COILS"
INPUT_REGISTERS = "INPUT_REGISTERS"
HOLDING_REGISTERS = "HOLDING_REGISTERS"
PRIMARY_TABLE = "primaryTable"
STARTING_ADDRESS = "startingAddress"
IS_BYTE_SWAP = "isByteSwap"
IS_WORD_SWAP = "isWordSwap"
// RAW_TYPE 定义了从Modbus设备上读取的二进制数据类型
RAW_TYPE = "rawType"
// STRING_REGISTER_SIZE 例如 "abcd" 需要 4 个字节也就需要使用两个寄存器(2个字), 所以 STRING_REGISTER_SIZE=2
STRING_REGISTER_SIZE = "stringRegisterSize"
SERVICE_STOP_WAIT_TIME = 1
)
var PrimaryTableBitCountMap = map[string]uint16{// 定义了各种寄存器对应的比特数
DISCRETES_INPUT: 1,
COILS: 1,
INPUT_REGISTERS: 16,
HOLDING_REGISTERS: 16,
}
var ValueTypeBitCountMap = map[string]uint16{// 定义了各种数据类型对应的比特数
common.ValueTypeInt16: 16,
common.ValueTypeInt32: 32,
common.ValueTypeInt64: 64,
common.ValueTypeUint16: 16,
common.ValueTypeUint32: 32,
common.ValueTypeUint64: 64,
common.ValueTypeFloat32: 32,
common.ValueTypeFloat64: 64,
common.ValueTypeBool: 1,
common.ValueTypeString: 16,
}
deviceclient.go
接口:
- 打开连接
- 关闭连接
- 获取值
- 设置值
结构体:
- Modbus生成指令所需的参数
- 对应DeviceProfile中的参数
功能:
- 创建指令信息
- 检查参数格式,错误则进行错误提示
- 对获取的参数进行格式调整
- 初始化指令信息结构体
- 计算寄存器长度
- 根据对应数据格式和寄存器的数据位数计算得到寄存器所需数量
- 二进制数据转换为实际结果
- 根据配置中数据类型进行转换
- 如果是可以交换的数据类型则根据配置决定是否交换
- 如果存在rawType则除了转换二进制到rawType,再转换到实际类型
- 实际输入转换为二进制数据
- 根据配置中数据类型进行转换,但注意不能超过对应寄存器数据长度
- 如果是可以交换的数据类型则根据配置决定是否交换
- 如果存在rawType则除了转换实际类型到rawType,再转换到二进制
- 计算寄存器对应字节数
- 转换二进制
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2018-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
package driver
import (
"bytes"
"encoding/binary"
"fmt"
"math"
"strings"
"time"
"github.com/edgexfoundry/device-sdk-go/v2/pkg/models"
"github.com/edgexfoundry/go-mod-core-contracts/v2/common"
"github.com/edgexfoundry/go-mod-core-contracts/v2/errors"
)
// DeviceClient 是一个需要 modbus client 库去实现的接口
// 它的职责是处理连接、读取数据字节值和写入数据字节值
type DeviceClient interface {
OpenConnection() error
GetValue(commandInfo interface{}) ([]byte, error)
SetValue(commandInfo interface{}, value []byte) error
CloseConnection() error
}
// CommandInfo 是命令信息
type CommandInfo struct {
PrimaryTable string
StartingAddress uint16
ValueType string
// 需要读取多少个寄存器
Length uint16
IsByteSwap bool
IsWordSwap bool
RawType string
}
func createCommandInfo(req *models.CommandRequest) (*CommandInfo, error) {// device profile参数验证,格式调整,最后创建命令信息结构体
if _, ok := req.Attributes[PRIMARY_TABLE]; !ok {
return nil, errors.NewCommonEdgeX(errors.KindContractInvalid, fmt.Sprintf("attribute %s not exists", PRIMARY_TABLE), nil)
}
primaryTable := fmt.Sprintf("%v", req.Attributes[PRIMARY_TABLE])
primaryTable = strings.ToUpper(primaryTable)
if _, ok := req.Attributes[STARTING_ADDRESS]; !ok {
return nil, errors.NewCommonEdgeX(errors.KindContractInvalid, fmt.Sprintf("attribute %s not exists", STARTING_ADDRESS), nil)
}
startingAddress, err := castStartingAddress(req.Attributes[STARTING_ADDRESS])
if err != nil {
return nil, errors.NewCommonEdgeX(errors.Kind(err), fmt.Sprintf("fail to cast %s", STARTING_ADDRESS), err)
}
var rawType = req.Type
if _, ok := req.Attributes[RAW_TYPE]; ok {
rawType = fmt.Sprintf("%v", req.Attributes[RAW_TYPE])
rawType, err = normalizeRawType(rawType)
if err != nil {
return nil, err
}
}
var length uint16
if req.Type == common.ValueTypeString {
length, err = castStartingAddress(req.Attributes[STRING_REGISTER_SIZE])
if err != nil {
return nil, err
} else if (length > 123) || (length < 1) {
return nil, errors.NewCommonEdgeX(errors.KindLimitExceeded, fmt.Sprintf("register size should be within the range of 1~123, get %v.", length), nil)
}
} else {
length = calculateAddressLength(primaryTable, rawType)
}
var isByteSwap = false
if _, ok := req.Attributes[IS_BYTE_SWAP]; ok {
isByteSwap, err = castSwapAttribute(req.Attributes[IS_BYTE_SWAP])
if err != nil {
return nil, errors.NewCommonEdgeX(errors.Kind(err), fmt.Sprintf("fail to cast %s", IS_BYTE_SWAP), err)
}
}
var isWordSwap = false
if _, ok := req.Attributes[IS_WORD_SWAP]; ok {
isWordSwap, err = castSwapAttribute(req.Attributes[IS_WORD_SWAP])
if err != nil {
return nil, errors.NewCommonEdgeX(errors.Kind(err), fmt.Sprintf("fail to cast %s", IS_WORD_SWAP), err)
}
}
return &CommandInfo{
PrimaryTable: primaryTable,
StartingAddress: startingAddress,
ValueType: req.Type,
Length: length,
IsByteSwap: isByteSwap,
IsWordSwap: isWordSwap,
RawType: rawType,
}, nil
}
func calculateAddressLength(primaryTable string, valueType string) uint16 { //根据对应数据格式和寄存器的数据位数计算得到寄存器所需数量
var primaryTableBit = PrimaryTableBitCountMap[primaryTable]
var valueTypeBitCount = ValueTypeBitCountMap[valueType]
var length = valueTypeBitCount / primaryTableBit
if length < 1 {
length = 1
}
return length
}
// TransformDataBytesToResult 用于将设备的二进制数据转换为指定的值类型作为实际结果
func TransformDataBytesToResult(req *models.CommandRequest, dataBytes []byte, commandInfo *CommandInfo) (*models.CommandValue, error) {
var err error
var res interface{}
var result = &models.CommandValue{}
switch commandInfo.ValueType {
case common.ValueTypeUint16:
res = binary.BigEndian.Uint16(dataBytes)
case common.ValueTypeUint32:
res = binary.BigEndian.Uint32(swap32BitDataBytes(dataBytes, commandInfo.IsByteSwap, commandInfo.IsWordSwap))
case common.ValueTypeUint64:
res = binary.BigEndian.Uint64(dataBytes)
case common.ValueTypeInt16:
res = int16(binary.BigEndian.Uint16(dataBytes))
case common.ValueTypeInt32:
res = int32(binary.BigEndian.Uint32(swap32BitDataBytes(dataBytes, commandInfo.IsByteSwap, commandInfo.IsWordSwap)))
case common.ValueTypeInt64:
res = int64(binary.BigEndian.Uint64(dataBytes))
case common.ValueTypeFloat32:
switch commandInfo.RawType {
case common.ValueTypeFloat32:
raw := binary.BigEndian.Uint32(swap32BitDataBytes(dataBytes, commandInfo.IsByteSwap, commandInfo.IsWordSwap))
res = math.Float32frombits(raw)
case common.ValueTypeInt16:
raw := int16(binary.BigEndian.Uint16(dataBytes))
res = float32(raw)
driver.Logger.Debugf("According to the rawType %s and the value type %s, convert integer %d to float %v ", INT16, FLOAT32, res, result.ValueToString())
case common.ValueTypeUint16:
raw := binary.BigEndian.Uint16(dataBytes)
res = float32(raw)
driver.Logger.Debugf("According to the rawType %s and the value type %s, convert integer %d to float %v ", UINT16, FLOAT32, res, result.ValueToString())
}
case common.ValueTypeFloat64:
switch commandInfo.RawType {
case common.ValueTypeFloat64:
raw := binary.BigEndian.Uint64(dataBytes)
res = math.Float64frombits(raw)
case common.ValueTypeInt16:
raw := int16(binary.BigEndian.Uint16(dataBytes))
res = float64(raw)
driver.Logger.Debugf("According to the rawType %s and the value type %s, convert integer %d to float %v ", INT16, FLOAT64, res, result.ValueToString())
case common.ValueTypeUint16:
raw := binary.BigEndian.Uint16(dataBytes)
res = float64(raw)
driver.Logger.Debugf("According to the rawType %s and the value type %s, convert integer %d to float %v ", UINT16, FLOAT64, res, result.ValueToString())
}
case common.ValueTypeBool:
res = false
// to find the 1st bit of the dataBytes by mask it with 2^0 = 1 (00000001)
if (dataBytes[0] & 1) > 0 {
res = true
}
case common.ValueTypeString:
res = string(bytes.Trim(dataBytes, string(rune(0))))
default:
return nil, fmt.Errorf("return result fail, none supported value type: %v", commandInfo.ValueType)
}
result, err = models.NewCommandValue(req.DeviceResourceName, commandInfo.ValueType, res)
if err != nil {
return nil, err
}
result.Origin = time.Now().UnixNano()
driver.Logger.Debugf("Transfer dataBytes to CommandValue(%v) successful.", result.ValueToString())
return result, nil
}
// TransformCommandValueToDataBytes 将读到的值转换为二进制数据,用于通过Modbus协议传输数据。
func TransformCommandValueToDataBytes(commandInfo *CommandInfo, value *models.CommandValue) ([]byte, error) {
var err error
var byteCount = calculateByteCount(commandInfo)
var dataBytes []byte
buf := new(bytes.Buffer)
if commandInfo.ValueType != common.ValueTypeString {// 转换非String类型数据
err = binary.Write(buf, binary.BigEndian, value.Value)
if err != nil {
return nil, fmt.Errorf("failed to transform %v to []byte", value.Value)
}
numericValue := buf.Bytes()
var maxSize = uint16(len(numericValue))
dataBytes = numericValue[maxSize-byteCount : maxSize]// 只取最后有数据的那一部分
}
_, ok := ValueTypeBitCountMap[commandInfo.ValueType]// 如果不支持该数据格式有错误提示
if !ok {
err = fmt.Errorf("none supported value type : %v \n", commandInfo.ValueType)
return dataBytes, err
}
if commandInfo.ValueType == common.ValueTypeUint32 || commandInfo.ValueType == common.ValueTypeInt32 || commandInfo.ValueType == common.ValueTypeFloat32 {// 可以进行swap操作的数据类型根据传入的是否交换的参数进行交换
dataBytes = swap32BitDataBytes(dataBytes, commandInfo.IsByteSwap, commandInfo.IsWordSwap)
}
// 根据rawType转换值,此功能将浮点值转换为32位整数值
if commandInfo.ValueType == common.ValueTypeFloat32 {
val, edgexErr := value.Float32Value()
if edgexErr != nil {
return dataBytes, edgexErr
}
if commandInfo.RawType == common.ValueTypeInt16 {
dataBytes, err = getBinaryData(int16(val))
if err != nil {
return dataBytes, err
}
} else if commandInfo.RawType == common.ValueTypeUint16 {
dataBytes, err = getBinaryData(uint16(val))
if err != nil {
return dataBytes, err
}
}
} else if commandInfo.ValueType == common.ValueTypeFloat64 {
val, edgexErr := value.Float64Value()
if edgexErr != nil {
return dataBytes, edgexErr
}
if commandInfo.RawType == common.ValueTypeInt16 {
dataBytes, err = getBinaryData(int16(val))
if err != nil {
return dataBytes, err
}
} else if commandInfo.RawType == common.ValueTypeUint16 {
dataBytes, err = getBinaryData(uint16(val))
if err != nil {
return dataBytes, err
}
}
} else if commandInfo.ValueType == common.ValueTypeString {// 转换String数据时不超过实际寄存器长度
// Cast value of string type
oriStr := value.ValueToString()
tempBytes := []byte(oriStr)
bytesL := len(tempBytes)
oriByteL := int(commandInfo.Length * 2)
if bytesL < oriByteL {
less := make([]byte, oriByteL-bytesL)
dataBytes = append(tempBytes, less...)
} else if bytesL > oriByteL {
dataBytes = tempBytes[:oriByteL]
} else {
dataBytes = []byte(oriStr)
}
}
driver.Logger.Debugf("Transfer CommandValue to dataBytes for write command, %v, %v", commandInfo.ValueType, dataBytes)
return dataBytes, err
}
func calculateByteCount(commandInfo *CommandInfo) uint16 {// 根据寄存器数量和寄存器类型计算得出数据的字节数
var byteCount uint16
if commandInfo.PrimaryTable == HOLDING_REGISTERS || commandInfo.PrimaryTable == INPUT_REGISTERS {
byteCount = commandInfo.Length * 2
} else {
byteCount = commandInfo.Length
}
return byteCount
}
func getBinaryData(val interface{}) (dataBytes []byte, err error) {// 转换成为二进制数
buf := new(bytes.Buffer)
err = binary.Write(buf, binary.BigEndian, val)// 将值用大端格式写入buf
if err != nil {
return dataBytes, err
}
dataBytes = buf.Bytes()// 转换成一个byte数组
return dataBytes, err
}
driver.go
流程:
- 处理读指令:
- 创建连接信息
- 创建一个设备客户端
- 结束后断开连接
- 处理命令请求
- 创建命令信息
- 从设备客户端取值
- 转换二进制数据为可读数据
- 处理写指令:
- 创建连接信息
- 创建一个设备客户端
- 结束后断开连接
- 处理命令请求
- 创建命令信息
- 转换读到的设置数据为二进制数据
- 往设备客户端写值
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2018-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
// Package driver is used to execute device-sdk's commands
package driver
import (
"fmt"
"sync"
"time"
sdkModel "github.com/edgexfoundry/device-sdk-go/v2/pkg/models"
"github.com/edgexfoundry/go-mod-core-contracts/v2/clients/logger"
"github.com/edgexfoundry/go-mod-core-contracts/v2/models"
)
var once sync.Once
var driver *Driver
type Driver struct {
Logger logger.LoggingClient
AsyncCh chan<- *sdkModel.AsyncValues
mutex sync.Mutex
addressMap map[string]chan bool
workingAddressCount map[string]int
stopped bool
}
var concurrentCommandLimit = 100
func (d *Driver) DisconnectDevice(deviceName string, protocols map[string]models.ProtocolProperties) error {
d.Logger.Warn("Driver's DisconnectDevice function didn't implement")
return nil
}
// lockAddress 所标记地址不可用,因为实际设备一次只能处理一个请求
func (d *Driver) lockAddress(address string) error {
if d.stopped {
return fmt.Errorf("service attempts to stop and unable to handle new request")
}
d.mutex.Lock()
lock, ok := d.addressMap[address]
if !ok {
lock = make(chan bool, 1)
d.addressMap[address] = lock
}
// workingAddressCount 用于检查高频命令执行情况,以避免goroutine阻塞
count, ok := d.workingAddressCount[address]
if !ok {
d.workingAddressCount[address] = 1
} else if count >= concurrentCommandLimit {
d.mutex.Unlock()
errorMessage := fmt.Sprintf("High-frequency command execution. There are %v commands with the same address in the queue", concurrentCommandLimit)
d.Logger.Error(errorMessage)
return fmt.Errorf(errorMessage)
} else {
d.workingAddressCount[address] = count + 1
}
d.mutex.Unlock()
lock <- true
return nil
}
// unlockAddress 在结束之后移除标记
func (d *Driver) unlockAddress(address string) {
d.mutex.Lock()
lock := d.addressMap[address]
d.workingAddressCount[address] = d.workingAddressCount[address] - 1
d.mutex.Unlock()
<-lock
}
// lockableAddress 根据协议返回可锁定的地址
func (d *Driver) lockableAddress(info *ConnectionInfo) string {
var address string
if info.Protocol == ProtocolTCP {
address = fmt.Sprintf("%s:%d", info.Address, info.Port)
} else {
address = info.Address
}
return address
}
func (d *Driver) HandleReadCommands(deviceName string, protocols map[string]models.ProtocolProperties, reqs []sdkModel.CommandRequest) (responses []*sdkModel.CommandValue, err error) {// 处理读指令
connectionInfo, err := createConnectionInfo(protocols)//创建连接信息
if err != nil {
driver.Logger.Errorf("Fail to create read command connection info. err:%v \n", err)
return responses, err
}
err = d.lockAddress(d.lockableAddress(connectionInfo))
if err != nil {
return responses, err
}
defer d.unlockAddress(d.lockableAddress(connectionInfo))
responses = make([]*sdkModel.CommandValue, len(reqs))
var deviceClient DeviceClient
// create device client and open connection
deviceClient, err = NewDeviceClient(connectionInfo)// 创建一个设备客户端
if err != nil {
driver.Logger.Infof("Read command NewDeviceClient failed. err:%v \n", err)
return responses, err
}
err = deviceClient.OpenConnection()// 连接设备
if err != nil {
driver.Logger.Infof("Read command OpenConnection failed. err:%v \n", err)
return responses, err
}
defer deviceClient.CloseConnection()// 结束后关闭设备
// 处理命令请求
for i, req := range reqs {
res, err := handleReadCommandRequest(deviceClient, req)
if err != nil {
driver.Logger.Infof("Read command failed. Cmd:%v err:%v \n", req.DeviceResourceName, err)
return responses, err
}
responses[i] = res
}
return responses, nil
}
func handleReadCommandRequest(deviceClient DeviceClient, req sdkModel.CommandRequest) (*sdkModel.CommandValue, error) {// 处理命令请求
var response []byte
var result = &sdkModel.CommandValue{}
var err error
commandInfo, err := createCommandInfo(&req)// 创建命令信息
if err != nil {
return nil, err
}
response, err = deviceClient.GetValue(commandInfo)// 从设备客户端取值
if err != nil {
return result, err
}
result, err = TransformDataBytesToResult(&req, response, commandInfo)// 转换二进制数据成可读数据
if err != nil {
return result, err
} else {
driver.Logger.Infof("Read command finished. Cmd:%v, %v \n", req.DeviceResourceName, result)
}
return result, nil
}
func (d *Driver) HandleWriteCommands(deviceName string, protocols map[string]models.ProtocolProperties, reqs []sdkModel.CommandRequest, params []*sdkModel.CommandValue) error {// 处理写指令
connectionInfo, err := createConnectionInfo(protocols)
if err != nil {
driver.Logger.Errorf("Fail to create write command connection info. err:%v \n", err)
return err
}
err = d.lockAddress(d.lockableAddress(connectionInfo))
if err != nil {
return err
}
defer d.unlockAddress(d.lockableAddress(connectionInfo))
var deviceClient DeviceClient
// create device client and open connection
deviceClient, err = NewDeviceClient(connectionInfo)// 新建设备客户端
if err != nil {
return err
}
err = deviceClient.OpenConnection()// 连接设备
if err != nil {
return err
}
defer deviceClient.CloseConnection()// 结束后断开连接
// 处理命令请求
for i, req := range reqs {
err = handleWriteCommandRequest(deviceClient, req, params[i])
if err != nil {
d.Logger.Error(err.Error())
break
}
}
return err
}
func handleWriteCommandRequest(deviceClient DeviceClient, req sdkModel.CommandRequest, param *sdkModel.CommandValue) error {// 处理写指令请求
var err error
commandInfo, err := createCommandInfo(&req)// 创建命令信息
if err != nil {
return err
}
dataBytes, err := TransformCommandValueToDataBytes(commandInfo, param)// 将写入数据转换成二进制数据便于发送
if err != nil {
return fmt.Errorf("transform command value failed, err: %v", err)
}
err = deviceClient.SetValue(commandInfo, dataBytes)// 往客户端中写入
if err != nil {
return fmt.Errorf("handle write command request failed, err: %v", err)
}
driver.Logger.Infof("Write command finished. Cmd:%v \n", req.DeviceResourceName)
return nil
}
func (d *Driver) Initialize(lc logger.LoggingClient, asyncCh chan<- *sdkModel.AsyncValues, deviceCh chan<- []sdkModel.DiscoveredDevice) error {
d.Logger = lc
d.AsyncCh = asyncCh
d.addressMap = make(map[string]chan bool)
d.workingAddressCount = make(map[string]int)
return nil
}
func (d *Driver) Stop(force bool) error {
d.stopped = true
if !force {
d.waitAllCommandsToFinish()
}
for _, locked := range d.addressMap {
close(locked)
}
return nil
}
// waitAllCommandsToFinish 用于检查和等待未完成的工作
func (d *Driver) waitAllCommandsToFinish() {
loop:
for {
for _, count := range d.workingAddressCount {
if count != 0 {
// wait a moment and check again
time.Sleep(time.Second * SERVICE_STOP_WAIT_TIME)
continue loop
}
}
break loop
}
}
func (d *Driver) AddDevice(deviceName string, protocols map[string]models.ProtocolProperties, adminState models.AdminState) error {
d.Logger.Debugf("Device %s is added", deviceName)
return nil
}
func (d *Driver) UpdateDevice(deviceName string, protocols map[string]models.ProtocolProperties, adminState models.AdminState) error {
d.Logger.Debugf("Device %s is updated", deviceName)
return nil
}
func (d *Driver) RemoveDevice(deviceName string, protocols map[string]models.ProtocolProperties) error {
d.Logger.Debugf("Device %s is removed", deviceName)
return nil
}
func NewProtocolDriver() sdkModel.ProtocolDriver {
once.Do(func() {
driver = new(Driver)
})
return driver
}
modbusclient.go
// -*- Mode: Go; indent-tabs-mode: t -*-
//
// Copyright (C) 2018-2021 IOTech Ltd
//
// SPDX-License-Identifier: Apache-2.0
package driver
import (
"encoding/binary"
"fmt"
"log"
"os"
"strings"
"time"
MODBUS "github.com/goburrow/modbus"
)
// ModbusClient 用于连接设备和读/写值
type ModbusClient struct {
// IsModbusTcp is a value indicating the connection type
IsModbusTcp bool
// TCPClientHandler is ued for holding device TCP connection
TCPClientHandler MODBUS.TCPClientHandler
// TCPClientHandler is ued for holding device RTU connection
RTUClientHandler MODBUS.RTUClientHandler
client MODBUS.Client
}
func (c *ModbusClient) OpenConnection() error {// 实现与Modbus设备连接
var err error
var newClient MODBUS.Client
if c.IsModbusTcp {
err = c.TCPClientHandler.Connect()
newClient = MODBUS.NewClient(&c.TCPClientHandler)
driver.Logger.Info(fmt.Sprintf("Modbus client create TCP connection."))
} else {
err = c.RTUClientHandler.Connect()
newClient = MODBUS.NewClient(&c.RTUClientHandler)
driver.Logger.Info(fmt.Sprintf("Modbus client create RTU connection."))
}
c.client = newClient
return err
}
func (c *ModbusClient) CloseConnection() error {// 实现与Modbus设备断开连接
var err error
if c.IsModbusTcp {
err = c.TCPClientHandler.Close()
} else {
err = c.RTUClientHandler.Close()
}
return err
}
func (c *ModbusClient) GetValue(commandInfo interface{}) ([]byte, error) {// 实现从设备中读取
var modbusCommandInfo = commandInfo.(*CommandInfo)
var response []byte
var err error
switch modbusCommandInfo.PrimaryTable {// 根据指令信息中的寄存器类型选择对应函数实现读指令
case DISCRETES_INPUT:
response, err = c.client.ReadDiscreteInputs(modbusCommandInfo.StartingAddress, modbusCommandInfo.Length)
case COILS:
response, err = c.client.ReadCoils(modbusCommandInfo.StartingAddress, modbusCommandInfo.Length)
case INPUT_REGISTERS:
response, err = c.client.ReadInputRegisters(modbusCommandInfo.StartingAddress, modbusCommandInfo.Length)
case HOLDING_REGISTERS:
response, err = c.client.ReadHoldingRegisters(modbusCommandInfo.StartingAddress, modbusCommandInfo.Length)
default:
driver.Logger.Error("None supported primary table! ")
}
if err != nil {
return response, err
}
driver.Logger.Info(fmt.Sprintf("Modbus client GetValue's results %v", response))
return response, nil
}
func (c *ModbusClient) SetValue(commandInfo interface{}, value []byte) error {//从设备中取值
var modbusCommandInfo = commandInfo.(*CommandInfo)
// Write value to device
var result []byte
var err error
switch modbusCommandInfo.PrimaryTable {// 根据指令信息中的寄存器类型选择对应函数实现写指令
case DISCRETES_INPUT:
err = fmt.Errorf("Error: DISCRETES_INPUT is Read-Only..!!")
case COILS:
result, err = c.client.WriteMultipleCoils(uint16(modbusCommandInfo.StartingAddress), modbusCommandInfo.Length, value)
case INPUT_REGISTERS:
err = fmt.Errorf("Error: INPUT_REGISTERS is Read-Only..!!")
case HOLDING_REGISTERS:
if modbusCommandInfo.Length == 1 {
result, err = c.client.WriteSingleRegister(uint16(modbusCommandInfo.StartingAddress), binary.BigEndian.Uint16(value))
} else {
result, err = c.client.WriteMultipleRegisters(uint16(modbusCommandInfo.StartingAddress), modbusCommandInfo.Length, value)
}
default:
}
if err != nil {
return err
}
driver.Logger.Info(fmt.Sprintf("Modbus client SetValue successful, results: %v", result))
return nil
}
func NewDeviceClient(connectionInfo *ConnectionInfo) (*ModbusClient, error) {// 将连接信息中的参数配置到ClientHandler中
client := new(ModbusClient)
var err error
if connectionInfo.Protocol == ProtocolTCP {
client.IsModbusTcp = true
}
if client.IsModbusTcp {
client.TCPClientHandler.Address = fmt.Sprintf("%s:%d", connectionInfo.Address, connectionInfo.Port)
client.TCPClientHandler.SlaveId = byte(connectionInfo.UnitID)
client.TCPClientHandler.Timeout = time.Duration(connectionInfo.Timeout) * time.Second
client.TCPClientHandler.IdleTimeout = time.Duration(connectionInfo.IdleTimeout) * time.Second
client.TCPClientHandler.Logger = log.New(os.Stdout, "", log.LstdFlags)
} else {
serialParams := strings.Split(connectionInfo.Address, ",")
client.RTUClientHandler.Address = serialParams[0]
client.RTUClientHandler.SlaveId = byte(connectionInfo.UnitID)
client.RTUClientHandler.Timeout = time.Duration(connectionInfo.Timeout) * time.Second
client.RTUClientHandler.IdleTimeout = time.Duration(connectionInfo.IdleTimeout) * time.Second
client.RTUClientHandler.BaudRate = connectionInfo.BaudRate
client.RTUClientHandler.DataBits = connectionInfo.DataBits
client.RTUClientHandler.StopBits = connectionInfo.StopBits
client.RTUClientHandler.Parity = connectionInfo.Parity
client.RTUClientHandler.Logger = log.New(os.Stdout, "", log.LstdFlags)
}
return client, err
}

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。
更多推荐


所有评论(0)