bin

edgex-launch.sh

目的

  • 启动所有EdgeX Go的二进制文件(必须以前就构建好)

前提

  • Consul和MongoDB都已经安装并运行

清理

  • 杀死进程edgex-device-modbus

流程

  1. 先到CMD目录下(…/cmd)
  2. 执行目录下的device-modbus文件,让其作为edgex-device-modbus进程启动
  3. 回到PWD目录
  4. 如果陷入trap,就执行清理功能并退出
  5. 循环运行监视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文件夹

流程

  1. 进入GIT_ROOT
  2. 检查有没有vendor.bk(备份),如果有以1退出,并提示该文件存在,需要先移除再继续
  3. 如果陷入trap执行清理功能
  4. 检查有没有vendor.bk,如果有将vendor.bk改为vendor(相当于备份还原)
  5. 如果vendor文件夹存在,备份他,这样我们就能建一个新的
  6. go mod时在项目中创建vendor文件将依赖包拷贝过来
  7. 打开nullglobbing,这样如果cmd目录中没有任何东西,那么我们就不会在这个循环中做任何事情
  8. 如果找不到Attribution.txt文件,以1退出,并提示该文件找不到
  9. 循环遍历vendor文件夹中的每个lib文件,确保lib存在且明确,否则提示文件中某lib丢失,请添加
  10. 回到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

流程

  1. 设置服务名称
  2. 建立对接的协议驱动
  3. 和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参数的结构体

功能

  • 创建连接信息
    1. 传入配置协议集
    2. 检查modbus协议形式
    3. 对错误配置进行错误提示
    4. 协议是RTU则调用功能创建RTU类型的连接信息
    5. 协议是TCP则调用功能创建TCP类型的连接信息
  • 解析整数配置参数
    1. 获取配置参数
    2. 转换成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中的参数

功能

  • 创建指令信息
    1. 检查参数格式,错误则进行错误提示
    2. 对获取的参数进行格式调整
    3. 初始化指令信息结构体
  • 计算寄存器长度
    • 根据对应数据格式和寄存器的数据位数计算得到寄存器所需数量
  • 二进制数据转换为实际结果
    1. 根据配置中数据类型进行转换
    2. 如果是可以交换的数据类型则根据配置决定是否交换
    3. 如果存在rawType则除了转换二进制到rawType,再转换到实际类型
  • 实际输入转换为二进制数据
    1. 根据配置中数据类型进行转换,但注意不能超过对应寄存器数据长度
    2. 如果是可以交换的数据类型则根据配置决定是否交换
    3. 如果存在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

流程

  • 处理读指令:
  1. 创建连接信息
  2. 创建一个设备客户端
  3. 结束后断开连接
  4. 处理命令请求
  5. 创建命令信息
  6. 从设备客户端取值
  7. 转换二进制数据为可读数据
  • 处理写指令:
  1. 创建连接信息
  2. 创建一个设备客户端
  3. 结束后断开连接
  4. 处理命令请求
  5. 创建命令信息
  6. 转换读到的设置数据为二进制数据
  7. 往设备客户端写值
// -*- 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
}

请添加图片描述

Logo

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

更多推荐