require "log" require "sys" require "utils" require "patch" require "pack" siganl = require 'posix.signal' syswait = require 'posix.sys.wait' unistd = require 'posix.unistd' mqtt=require "mosquitto" base64=require("base64") dyutils=require("dyutils") BATT_SIZE = 18 RACK_PATH = "/app/rack/" RACK_CONFIG_FILE = RACK_PATH .. "threshold.cfg" RACK_STATE_FILE = RACK_PATH .. "alarm.dat" LOG_LEVEL = log.LOGLEVEL_TRACE -- 函数重载 local bpack = function (format, ... ) if not format then return nil end local dy_format = format dy_format = string.gsub(dy_format,'C','b') dy_format = string.gsub(dy_format,'s','h') dy_format = string.gsub(dy_format,'S','H') return bpack( dy_format, ... ) end local bunpack = function ( string, format, init ) if not string or not format then return nil end local dy_format = format dy_format = string.gsub(dy_format,'C','b') dy_format = string.gsub(dy_format,'s','h') dy_format = string.gsub(dy_format,'S','H') return bunpack( string, dy_format, init ) end nodes={} gw_sn=nil template_id=nil recvBuff={} alarm_threshold = {RackUnderVolAlarm={100,120,"Under"}} alarm_state = {RackUnderVolAlarm=0} ----------------------[MODBUS read data]-------------------------- local function read_multi_register(obj,start,bytes_len,cmd) return cmd or 0x03,start,bytes_len/2 end -- ----------------------[MODBUS write data]-------------------------- local function write_multi_register(obj,start,data,bytes_len,cmd) return cmd or 0x10,start,bytes_len/2,bytes_len,data end ----------------------[MODBUS write data]-------------------------- local function write_once_register(obj,start,data,cmd) return cmd or 0x05,start,nil,nil,data end ----------------------[数据采集解析]-------------------------- local function response_RackCellTempQuery(raw_data, len) local obj={values={}} for i=0,BATT_SIZE-1 do local _, value= bunpack(raw_data,">S") if not value then return end raw_data = string.sub(raw_data,2 + 1 ,-1) table.insert( obj.values,value*0.01) end return obj end local function response_RackCellVolQuery(raw_data, len) local obj={values={}} for i=0,BATT_SIZE-1 do local _, value= bunpack(raw_data,">S") if not value then return end raw_data = string.sub(raw_data,2 + 1 ,-1) table.insert( obj.values,value*0.001) end return obj end local function response_RackCellIntResQuery(raw_data, len) local obj={values={}} for i=0,BATT_SIZE-1 do local _, value= bunpack(raw_data,">S") if not value then return end raw_data = string.sub(raw_data,2 + 1 ,-1) table.insert( obj.values,value*0.001) end return obj end local function response_RackCellStateQuery(raw_data, len) local obj={values={}} for i=0,BATT_SIZE-1 do local _, value= bunpack(raw_data,">S") if not value then return end raw_data = string.sub(raw_data,2 + 1 ,-1) local bits = { B1=bit.band(bit.rshift(value,1),1), B3=bit.band(bit.rshift(value,3),1), B4=bit.band(bit.rshift(value,4),1), B5=bit.band(bit.rshift(value,5),1), B6=bit.band(bit.rshift(value,6),1), B9=bit.band(bit.rshift(value,9),1), B15=bit.band(bit.rshift(value,15),1) } table.insert( obj.values,bits) end return obj end local function response_RackSystemInfo(raw_data, len) local obj={} local _, state= bunpack(raw_data,">S") if not state then return end raw_data = string.sub(raw_data,2 + 1 ,-1) local _, value= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["RackTotalVoltage"] = value*0.01 local _, value= bunpack(raw_data,">s") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["RackTotalCurrent"] = value*0.01 local _, value= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["RackSoc"] = value*0.1 local _, value= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["RackSoh"] = value*0.1 local _, value= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["CellCount"] = value local set,clr = 0,0 if obj["RackTotalVoltage"] < alarm_threshold.RackUnderVolAlarm[1] then set = 1 elseif obj["RackTotalVoltage"] >= alarm_threshold.RackUnderVolAlarm[2] then clr = 1 end if (alarm_state.RackUnderVolAlarm == 0 and set == 1 ) or (alarm_state.RackUnderVolAlarm == 1 and clr == 1 ) then log.info("alarm","state of set",set,"old",alarm_state.RackUnderVolAlarm) alarm_state.RackUnderVolAlarm = set -- Event(gw_sn,dev_sn,FetchPortBySN(dev_sn) or "NULL","RackAlarm",{AlarmFlag={B0=alarm_state.RackUnderVolAlarm}}) io.writeFile(RACK_STATE_FILE,cjson.encode(alarm_state)) end obj["AlarmFlag"] = { B0=bit.band(bit.rshift(state,0),1), B1=bit.band(bit.rshift(state,1),1), B2=bit.band(bit.rshift(state,2),1), B3=alarm_state.RackUnderVolAlarm, B4=bit.band(bit.rshift(state,4),1), -- B6=bit.band(bit.rshift(state,6),1), -- B7=bit.band(bit.rshift(state,7),1), -- B8=bit.band(bit.rshift(state,8),1), B9=bit.band(bit.rshift(state,9),1), B15=bit.band(bit.rshift(state,15),1) } obj["RunState"] = 0 if bit.band(bit.rshift(state,6),1) == 1 then obj["RunState"] = 1 elseif bit.band(bit.rshift(state,7),1) == 1 then obj["RunState"] = 2 elseif bit.band(bit.rshift(state,8),1) == 1 then obj["RunState"] = 3 end return obj end global_identifier_func_tab = { { identifier="RackCellTempQuery", create_bin_func=read_multi_register, create_bin_args={0x0409,BATT_SIZE*2}, parser_bin_func=response_RackCellTempQuery }, { identifier="RackCellVolQuery", create_bin_func=read_multi_register, create_bin_args={0x0309,BATT_SIZE*2}, parser_bin_func=response_RackCellVolQuery }, { identifier="RackCellIntResQuery", create_bin_func=read_multi_register, create_bin_args={0x0609,BATT_SIZE*2}, parser_bin_func=response_RackCellIntResQuery }, { identifier="RackCellStateQuery", create_bin_func=read_multi_register, create_bin_args={0x0509,BATT_SIZE*2}, parser_bin_func=response_RackCellStateQuery }, { identifier="RackSystemInfo", create_bin_func=read_multi_register, create_bin_args={0x0303,6*2}, parser_bin_func=response_RackSystemInfo } } -- 事件上报 --[[ local function Event(gw_sn,dev_sn,port,identifier,param,mi,time) local topic = string.format( "ipc/%s/%s/device/%s/data/event/%s",gw_sn,port,dev_sn,identifier) local obj_send={} local json_str obj_send['mi'] = mi or 0 -- obj_send['len'] = #data obj_send['identifier'] = identifier obj_send['sn'] = dev_sn obj_send['time'] = time or os.time() obj_send["data_type"] = 1 obj_send["tag_node"] = type(param) == "string" and param or cjson.encode(param) json_str = cjson.encode(obj_send) log.info("Event",topic,json_str) client:publish(topic, json_str) end --]] ------------------------------------- [串口 -> 平台] 解析输出 ------------------------------------- local function Port2MessageOutput(gw_sn,dev_sn,port,identifier,param, mi) if not gw_sn or not port or not dev_sn then log.error("Port2MessageOutput", "param error",gw_sn,dev_sn,port,identifier) return end topic = string.format( "ipc/%s/%s/device/%s/data/service/%s",gw_sn,port,dev_sn,identifier) obj_send={} obj_send['mi'] = mi -- obj_send['len'] = #data obj_send['identifier'] = identifier obj_send['sn'] = dev_sn obj_send['port'] = port obj_send['time'] = os.time() obj_send['disable_report'] = 0 obj_send['report_period'] = 0 obj_send['data_type'] = 2 obj_send['userParam'] = "" if param and type(param) == "table" then obj_send['tag_node']=cjson.encode(param) end json_str = cjson.encode(obj_send) log.info("Port2MessageOutput",topic,json_str) -- mqtt_client:publish(topic, json_str) client:publish(topic, json_str) end ------------------------------------- [平台 -> 串口] 解析输出 ------------------------------------- local function Message2PortOutput(gw_sn,dev_sn,port,identifier,data,period,mi) if not gw_sn or not port or not dev_sn then log.error("Message2PortOutput", "param error",gw_sn,dev_sn,port,identifier,data) return end topic = string.format( "ipc/%s/%s/device/%s/data/Set_Rglt_Raw",gw_sn,port,dev_sn) obj_send={} obj_send['mi'] = mi obj_send['len'] = #data obj_send['src_identifier'] = identifier obj_send['sn'] = dev_sn obj_send['port'] = port obj_send['protocol'] = 3 obj_send['communication_timeout'] = 200 obj_send['period'] = period or 0 obj_send['term_addr'] = SUB_ADDR obj_send['data_b64'] = base64.encode(data) json_str = cjson.encode(obj_send) log.info("Message2PortOutput",topic,json_str) -- mqtt_client:publish(topic, json_str) client:publish(topic, json_str) end local function FetchAddrBySN(sn) for _,v in pairs(nodes) do if v.sn == sn then return v.term_addr end end end local function FetchPortBySN(sn) for _,v in pairs(nodes) do if v.sn == sn then return v.connect_port end end end local function protocol_encode(json_str, len) local err_code=0 --取出入参-- local obj=cjson.decode(json_str) local id=obj.identifier local server_period=obj.server_period local dev_sn=obj.sn local addrs = FetchAddrBySN(dev_sn) if not addrs then return end local addr=tonumber(addrs,16) --编码-- if id == "null" then log.warn("protocol_encode","Need device sn in payload") return 0, "", 0 end log.info("protocol_encode","input:"..json_str) if obj.identifier == "period" then sys.publish("SCHEDULE_PERIOD",obj.sn,cjson.decode(obj.input)) return elseif obj.identifier == "oneshot" then sys.publish("SCHEDULE_ONESHOT",obj.sn,cjson.decode(obj.input)) return elseif obj.identifier == "FetchAlarmThreshold" then Port2MessageOutput(gw_sn,dev_sn,"NULL",obj.identifier,{RackUnderVolAlarm = alarm_threshold.RackUnderVolAlarm[1],RackUnderVolAlarmRecov = alarm_threshold.RackUnderVolAlarm[2]}, obj.mi) return elseif obj.identifier == "RackAlarmConfig" then if obj.tags and obj.tags.RackUnderVolAlarm and obj.tags.RackUnderVolAlarmRecov then alarm_threshold.RackUnderVolAlarm[1] = obj.tags.RackUnderVolAlarm alarm_threshold.RackUnderVolAlarm[2] = obj.tags.RackUnderVolAlarmRecov io.writeFile(RACK_CONFIG_FILE,cjson.encode(alarm_threshold)) log.info("config","alarm config update") Port2MessageOutput(gw_sn,dev_sn,"NULL",obj.identifier,{}, obj.mi) end return end local raw_data=nil for i,v in pairs(global_identifier_func_tab) do if v.identifier == id then local call_func = v.create_bin_func local func,start,len,data_len,data = call_func(obj,v.create_bin_args[1],v.create_bin_args[2],v.create_bin_args[3]) local bin=nil if not len then raw_data = bpack("CC>S",addr,func,start) for i = 1, #data do raw_data = raw_data .. bpack("c",data:byte(i)) end elseif not data then raw_data = bpack("CC>S>S",addr,func,start,len) else raw_data = bpack("CC>S>SC",addr,func,start,len,data_len) for i = 1, #data do raw_data = raw_data .. bpack("c",data:byte(i)) end end raw_data = raw_data .. bpack(">S",dyutils.CRC16(raw_data,#raw_data)) Message2PortOutput(gw_sn,dev_sn,FetchPortBySN(dev_sn),id,raw_data,server_period, obj.mi) return end end return -1, nil, 0 end local function protocol_decode(input_str, len) local err_code=0 --取出入参-- local obj={} local input_obj=cjson.decode(input_str) local id=input_obj.src_identifier local dev_sn=input_obj.sn --入参中的raw_data的二进制转换-- raw_data = base64.decode(input_obj.data_b64) _,addr,func = bunpack(raw_data,"CC") local addrs = FetchAddrBySN(dev_sn) if not addrs or addr~=tonumber(addrs,16) then return nil end local crc = dyutils.CRC16(raw_data,#raw_data - 2) _,crc_pack = bunpack(string.sub(raw_data,#raw_data - 1),">S") if (crc ~= crc_pack) then log.info(string.format( "Data crc check failed cal:%04X packet:%04X",crc,crc_pack)) return nil end for i,v in pairs(global_identifier_func_tab) do if v.identifier == id then local call_func = v.parser_bin_func local data,data_len = nil,0 if func == 0x05 then data_len = #raw_data - 5 - 1 data = string.sub(raw_data,5,-3) elseif func == 0x01 or func == 0x02 or func == 0x03 or func == 0x04 then data_len = #raw_data - 4 - 1 data = string.sub(raw_data,4,-3) elseif func == 0x10 then data_len = 0 data = nil else log.info(string.format( "unknow function code %02X",func)) return -1 end log.info("protocol_decode","protocol_decode raw_data ",string.toHex(data) ) local json_obj = call_func(data,data_len)-- 调用不同的函数执行解析数据 if not json_obj then -- log.info("decode failed,raw_data",string.toHex(data),data_len) log.warn("protocol_decode","decode failed,raw_data:",string.toHex(data)) return end -- json_obj.identifier=id -- 返回调用标识 -- json_obj.term_addr=string.format( "%02X",addr) -- json_obj.term_addr=addr Port2MessageOutput(gw_sn,dev_sn,FetchPortBySN(dev_sn) or "NULL",id,json_obj, input_obj.mi) return end end return nil end ---------------------------------------------------------- 外部消息---------------------------------------------------------- local function mqtt_session(id, addr, port, usr, pwd, subscribe, qos,func) cid, keepAlive, timeout = tonumber(cid) or 1, tonumber(keepAlive) or 300, tonumber(timeout) or 1800 cleansession, qos, retain = tonumber(cleansession) or 0, tonumber(qos) or 0, tonumber(retain) or 0 sys.taskInit(function() mqtt.init() client = mqtt.new() client.ON_CONNECT = function() log.warn("mqtt","client connected") for _,v in pairs(subscribe) do if type(v) == "string" then client:subscribe(v) end end end client.ON_MESSAGE = function(mid, topic, payload) log.info("mqtt recv",topic, payload) if topic then -- 控制命令通道直接调用 table.insert(recvBuff,{topic,payload or ''}) sys.publish("INTERNAL_CLIENT_RECV_INT") end end client.ON_DISCONNECT = function() log.warn("mqtt","client disconnected") end client:connect(addr,port) sys.poll_socket=function() client:loop(3000,1) end while(true) do while client:socket() do sys.wait(2000) end sys.wait(1000) client:reconnect() end end) end function on_mqtt_msg(topic, payload) -- ipc/21012C000038/VIRTUAL_1/device/333333333/data_filtered/property/vdev_data if string.find(topic,"data/Set_Rglt") then local ret,err = pcall(protocol_encode,payload,#payload) if not ret then log.info("mqtt.recv",err) end elseif string.find(topic,"data/raw_data") then --二进制数据流解析 local ret,err = pcall(protocol_decode,payload,#payload) if not ret then log.info("mqtt.recv",err) end else log.info("due message:",topic, payload) end end ------------------------------------- 上文是控制逻辑 ------------------------------------- function main(nodes_cfg,sn) log.info("main_loop", "=============== JING HU BMS RUN ============") if not nodes_cfg then return -1,"Input param error" end gw_sn = sn log.info("main_loop", string.format("Try Run At Shell\nsn='%s'\nnodes='%s'\nmain(nodes,sn)",gw_sn,nodes_cfg)) local obj = cjson.decode(nodes_cfg) -- if obj and obj.nodes_cfg then -- nodes = obj.nodes_cfgS -- if next(nodes) == nil then -- return -2,"not any node" -- end if obj and #obj then nodes = obj else return -3,"Input param miss word" end local config = io.readFile(RACK_CONFIG_FILE) if config then local obj = cjson.decode(config) if obj then alarm_threshold = obj end end local config = io.readFile(RACK_STATE_FILE) if config then local obj = cjson.decode(config) if obj then alarm_state = obj end end local subscribe={"ipc/+/+/device/+/data/raw_data","ipc/+/+/device/+/data/Set_Rglt"} mqtt_session("mandun","localhost", 1883 , nil, nil, subscribe, 0) -- mqtt received and process sys.taskInit(function() while true do sys.waitUntil("INTERNAL_CLIENT_RECV_INT",2000) if #recvBuff > 0 then local packet = table.remove(recvBuff) if packet then on_mqtt_msg(packet[1],packet[2]) packet = nil end end end end) --启动系统框架 sys.init(0, 0) sys.run() end main(arg[1],arg[2]) -- sn='218805000002' -- nodes='[{"connect_port":"RS485_1","depth":0,"product_key":"44521175","sn":"0i37jD9e","template_id":"44521175","term_addr":"1"}]' -- main(nodes,sn)