require "log" require "sys" require "utils" require "patch" lpack=require("lua_pack") MQTT=require("mqtt_library") base64=require("base64") dyutils=require("dyutils") nodes={} gw_sn=nil template_id=nil ----------------------[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 control_hezha(obj) local cmd = 0x00 if obj and obj.kongkaibianhao and obj.hezhama then local data = bpack(">S",obj.hezhama) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x05) end end local function response_hezha(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["hezhama"] = hezhama return obj end -----------设置温度预警门限 local function control_setwenduyujing(obj) local cmd = 0x0B if obj and obj.kongkaibianhao and obj.wenduyujingmenxian then local data = bpack(">S",obj.wenduyujingmenxian*10) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setwenduyujing(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置功率门限 local function control_setgonglumenxian(obj) local cmd = 0x03 if obj and obj.kongkaibianhao and obj.gonglumenxianzhi then local data = bpack(">S",obj.gonglumenxianzhi) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setgonglumenxian(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置漏电预警 local function control_setloudianyujing(obj) local cmd = 0x0A if obj and obj.kongkaibianhao and obj.loudianyujingzhi then local data = bpack(">S",obj.loudianyujingzhi*10) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setloudianyujing(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置A相电流预警门限 local function control_setAxiangdianliuyujing(obj) local cmd = 0x0C if obj and obj.kongkaibianhao and obj.Axiangdianliuyujingmenxian then local data = bpack(">S",obj.Axiangdianliuyujingmenxian*100) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setAxiangdianliuyujing(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置B相电流预警门限 local function control_setBxiangdianliuyujing(obj) local cmd = 0x0D if obj and obj.kongkaibianhao and obj.Bxiangdianliuyujingmenxian then local data = bpack(">S",obj.Bxiangdianliuyujingmenxian*100) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setBxiangdianliuyujing(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置C相电流预警门限 local function control_setCxiangdianliuyujing(obj) local cmd = 0x0E if obj and obj.kongkaibianhao and obj.Cxiangdianliuyujingmenxian then local data = bpack(">S",obj.Cxiangdianliuyujingmenxian*100) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setCxiangdianliuyujing(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置预警电压阈值上限 local function control_setyujingdianyayuzhishangxian(obj) local cmd = 0x14 if obj and obj.kongkaibianhao and obj.yujingdianyayuzhishangxian then local data = bpack(">S",obj.yujingdianyayuzhishangxian) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setyujingdianyayuzhishangxian(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end -----------设置预警电压阈值下限 local function control_setyujingdianyayuzhixiaxian(obj) local cmd = 0x15 if obj and obj.kongkaibianhao and obj.yujingdianyayuzhixiaxian then local data = bpack(">S",obj.yujingdianyayuzhixiaxian) return write_once_register(obj,bit.lshift(cmd,8) + obj.kongkaibianhao,data,0x06) end end local function response_setyujingdianyayuzhixiaxian(raw_data, len) local obj={} local _, hezhama= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["fanhuizhi"] = fanhuizhi return obj end ----------------------[数据采集解析]-------------------------- local function response_Heartbeat(raw_data, len) obj={} local _, Heartbeat= bunpack(raw_data,">S") obj["Heartbeat"] = 200 return obj end local function response_guigexinghao(raw_data, len) obj={} for i=0,29 do local _, guige,xinghao= bunpack(raw_data,"CC") if not guige or not xinghao then return obj end raw_data = string.sub(raw_data,2 + 1 ,-1) obj["guige".. i] = guige obj["xinghao".. i] = bit.band(xinghao,0xF) end return obj end local function response_duqukaihezha(raw_data, len) obj={} local format = nil local bits = nil if len == 1 then format = "C" elseif len == 2 then format = " 16) and 29 or len*8 - 1 do obj["kaihezhazhuangtai".. i] = bit.band(bit.rshift(kaihezhazhuangtai,i),1) end return obj end local function response_gaojing(raw_data, len) obj={} for i=0,29 do local _, gaojing= bunpack(raw_data,">S") if not gaojing then return obj end raw_data = string.sub(raw_data,2 + 1 ,-1) obj["langyongbaojing".. i] = bit.band(bit.rshift(gaojing,1),1) obj["guozaibaojing".. i] = bit.band(bit.rshift(gaojing,2),1) obj["wendubaojing".. i] = bit.band(bit.rshift(gaojing,3),1) obj["loudianbaojing".. i] = bit.band(bit.rshift(gaojing,4),1) obj["guoliubaojing".. i] = bit.band(bit.rshift(gaojing,5),1) obj["guoyabaojing".. i] = bit.band(bit.rshift(gaojing,6),1) obj["loudianbaohugongnengzhengchang".. i] = bit.band(bit.rshift(gaojing,7),1) obj["loudianbaohuzijianweiwancheng".. i] = bit.band(bit.rshift(gaojing,8),1) obj["shuruquexiangbaojing3P4P".. i] = bit.band(bit.rshift(gaojing,9),1) end return obj end local function response_xianludianliu(raw_data, len) obj={} for i=0,29 do local _, xianludianliu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["xianludianliu".. i] = xianludianliu*0.01 end return obj end local function response_xianludianya(raw_data, len) obj={} for i=0,29 do local _, xianludianya= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["xianludianya".. i] = xianludianya end return obj end --漏电电流 local function response_loudiandianliu(raw_data, len) obj={} for i=0,29 do local _, loudiandianliu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["loudiandianliu".. i] = loudiandianliu*0.1 end return obj end --A相电压 local function response_Axiangdianya(raw_data, len) obj={} for i=0,29 do local _, Axiangdianya= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Axiangdianya".. i] = Axiangdianya end return obj end --B相电压 local function response_Bxiangdianya(raw_data, len) obj={} for i=0,29 do local _, Bxiangdianya= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Bxiangdianya".. i] = Bxiangdianya end return obj end --C相电压 local function response_Cxiangdianya(raw_data, len) obj={} for i=0,29 do local _, Cxiangdianya= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Cxiangdianya".. i] = Cxiangdianya end return obj end --A相电流 local function response_Axiangdianliu(raw_data, len) obj={} for i=0,29 do local _, Axiangdianliu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Axiangdianliu".. i] = Axiangdianliu*0.01 end return obj end --B相电流 local function response_Bxiangdianliu(raw_data, len) obj={} for i=0,29 do local _, Bxiangdianliu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Bxiangdianliu".. i] = Bxiangdianliu*0.01 end return obj end --C相电流 local function response_Cxiangdianliu(raw_data, len) obj={} for i=0,29 do local _, Cxiangdianliu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Cxiangdianliu".. i] = Cxiangdianliu*0.01 end return obj end local function response_dugonglumenxian(raw_data, len) obj={} for i=0,29 do local _, xianlugonglumenxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["xianlugonglumenxian".. i] = xianlugonglumenxian end return obj end local function response_dudianyayujingshangxian(raw_data, len) obj={} for i=0,29 do local _, dianyayujingshangxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["dianyayujingshangxian".. i] = dianyayujingshangxian end return obj end local function response_dudianyayujingxiaxian(raw_data, len) obj={} for i=0,29 do local _, dianyayujingxiaxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["dianyayujingxiaxian".. i] = dianyayujingxiaxian end return obj end local function response_duloudianyujingshangxian(raw_data, len) obj={} for i=0,29 do local _, loudianyujingshangxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["loudianyujingshangxian".. i] = loudianyujingshangxian*0.1 end return obj end local function response_duyujingwendumenxian(raw_data, len) obj={} for i=0,29 do local _, yujingwendumenxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["yujingwendumenxian".. i] = yujingwendumenxian*0.1 end return obj end local function response_duAxiangdianliuyujing(raw_data, len) obj={} for i=0,29 do local _, Axiangdianliuyujing= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Axiangdianliuyujing".. i] = Axiangdianliuyujing*0.01 end return obj end local function response_duBxiangdianliuyujing(raw_data, len) obj={} for i=0,29 do local _, Bxiangdianliuyujing= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Bxiangdianliuyujing".. i] = Bxiangdianliuyujing*0.01 end return obj end local function response_duCxiangdianliuyujing(raw_data, len) obj={} for i=0,29 do local _, Cxiangdianliuyujing= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["Cxiangdianliuyujing".. i] = Cxiangdianliuyujing*0.01 end return obj end local function response_dudianliuyujingmenxian(raw_data, len) obj={} for i=0,29 do local _, dianliuyujingmenxian= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["dianliuyujingmenxian".. i] = dianliuyujingmenxian*0.01 end return obj end local function response_dianliangdi(raw_data, len) obj={} local _, dianliangdi= bunpack(raw_data,">S") obj["dianliangdi"] = dianliangdi return obj end local function response_dianlianggao(raw_data, len) obj={} local _, dianlianggao= bunpack(raw_data,">S") obj["dianlianggao"] = dianlianggao return obj end local function response_kongkaishu(raw_data, len) obj={} local format = nil local bits = nil if len == 1 then format = "C" elseif len == 2 then format = " 16) and 29 or len*8 - 1 do obj["kongkaishu".. i] = bit.band(bit.rshift(kongkaishu,i),1) end return obj end local function response_mokuaiwendu(raw_data, len) obj={} for i=0,29 do local _, mokuaiwendu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["mokuaiwendu".. i] = mokuaiwendu*0.1 end return obj end local function response_xianlugonglu(raw_data, len) obj={} for i=0,29 do local _, xianlugonglu= bunpack(raw_data,">S") raw_data = string.sub(raw_data,2 + 1 ,-1) obj["xianlugonglu".. i] = xianlugonglu end return obj end global_identifier_func_tab = { { identifier="Heartbeat", create_bin_func=read_multi_register, create_bin_args={0x0000,2,0x03}, parser_bin_func=response_Heartbeat }, { identifier="guigexinghao", create_bin_func=read_multi_register, create_bin_args={0x0600,60,0x04}, parser_bin_func=response_guigexinghao }, { identifier="duqukaihezha", create_bin_func=read_multi_register, create_bin_args={0x0000,60,0x01}, parser_bin_func=response_duqukaihezha }, { identifier="gaojing", create_bin_func=read_multi_register, create_bin_args={0x0500,60,0x03}, parser_bin_func=response_gaojing }, { identifier="xianludianliu", create_bin_func=read_multi_register, create_bin_args={0x0400,60,0x03}, parser_bin_func=response_xianludianliu }, { identifier="xianludianya", create_bin_func=read_multi_register, create_bin_args={0x0000,60,0x03}, parser_bin_func=response_xianludianya }, { identifier="dugonglumenxian", create_bin_func=read_multi_register, create_bin_args={0x0300,60,0x04}, parser_bin_func=response_dugonglumenxian }, { identifier="dudianyayujingshangxian", create_bin_func=read_multi_register, create_bin_args={0x0F00,60,0x04}, parser_bin_func=response_dudianyayujingshangxian }, { identifier="dudianyayujingxiaxian", create_bin_func=read_multi_register, create_bin_args={0x1000,60,0x04}, parser_bin_func=response_dudianyayujingxiaxian }, { identifier="duloudianyujingshangxian", create_bin_func=read_multi_register, create_bin_args={0x1100,60,0x04}, parser_bin_func=response_duloudianyujingshangxian }, { identifier="duyujingwendumenxian", create_bin_func=read_multi_register, create_bin_args={0x1200,60,0x04}, parser_bin_func=response_duyujingwendumenxian }, { identifier="duAxiangdianliuyujing", create_bin_func=read_multi_register, create_bin_args={0x1300,60,0x04}, parser_bin_func=response_duAxiangdianliuyujing }, { identifier="duBxiangdianliuyujing", create_bin_func=read_multi_register, create_bin_args={0x1400,60,0x04}, parser_bin_func=response_duBxiangdianliuyujing }, { identifier="duCxiangdianliuyujing", create_bin_func=read_multi_register, create_bin_args={0x1500,60,0x04}, parser_bin_func=response_duCxiangdianliuyujing }, { identifier="dudianliuyujingmenxian", create_bin_func=read_multi_register, create_bin_args={0x1600,60,0x04}, parser_bin_func=response_dudianliuyujingmenxian }, { identifier="dianliangdi", create_bin_func=read_multi_register, create_bin_args={0x0600,2,0x03}, parser_bin_func=response_dianliangdi }, { identifier="dianlianggao", create_bin_func=read_multi_register, create_bin_args={0x0700,2,0x03}, parser_bin_func=response_dianlianggao }, { identifier="kongkaishu", create_bin_func=read_multi_register, create_bin_args={0x0000,0xff*2,0x01}, parser_bin_func=response_kongkaishu }, { identifier="mokuaiwendu", create_bin_func=read_multi_register, create_bin_args={0x0300,60,0x03}, parser_bin_func=response_mokuaiwendu }, { identifier="xianlugonglu", create_bin_func=read_multi_register, create_bin_args={0x0200,60,0x03}, parser_bin_func=response_xianlugonglu }, { identifier="loudiandianliu", create_bin_func=read_multi_register, create_bin_args={0x0100,60,0x03}, parser_bin_func=response_loudiandianliu }, { identifier="Axiangdianya", create_bin_func=read_multi_register, create_bin_args={0x0800,60,0x03}, parser_bin_func=response_Axiangdianya }, { identifier="Bxiangdianya", create_bin_func=read_multi_register, create_bin_args={0x0900,60,0x03}, parser_bin_func=response_Bxiangdianya }, { identifier="Cxiangdianya", create_bin_func=read_multi_register, create_bin_args={0x0A00,60,0x03}, parser_bin_func=response_Cxiangdianya }, { identifier="Axiangdianliu", create_bin_func=read_multi_register, create_bin_args={0x0B00,60,0x03}, parser_bin_func=response_Axiangdianliu }, { identifier="Bxiangdianliu", create_bin_func=read_multi_register, create_bin_args={0x0C00,60,0x03}, parser_bin_func=response_Bxiangdianliu }, { identifier="Cxiangdianliu", create_bin_func=read_multi_register, create_bin_args={0x0D00,60,0x03}, parser_bin_func=response_Cxiangdianliu }, { identifier="hezha", create_bin_func=control_hezha, create_bin_args={nil,nil}, parser_bin_func=response_hezha }, { identifier="setwenduyujing", create_bin_func=control_setwenduyujing, create_bin_args={nil,nil}, parser_bin_func=response_setwenduyujing }, { identifier="setgonglumenxian", create_bin_func=control_setgonglumenxian, create_bin_args={nil,nil}, parser_bin_func=response_setgonglumenxian }, { identifier="setloudianyujing", create_bin_func=control_setloudianyujing, create_bin_args={nil,nil}, parser_bin_func=response_setloudianyujing }, { identifier="setAxiangdianliuyujing", create_bin_func=control_setAxiangdianliuyujing, create_bin_args={nil,nil}, parser_bin_func=response_setAxiangdianliuyujing }, { identifier="setBxiangdianliuyujing", create_bin_func=control_setBxiangdianliuyujing, create_bin_args={nil,nil}, parser_bin_func=response_setBxiangdianliuyujing }, { identifier="setCxiangdianliuyujing", create_bin_func=control_setCxiangdianliuyujing, create_bin_args={nil,nil}, parser_bin_func=response_setCxiangdianliuyujing }, { identifier="setyujingdianyayuzhishangxian", create_bin_func=control_setyujingdianyayuzhishangxian, create_bin_args={nil,nil}, parser_bin_func=response_setyujingdianyayuzhishangxian }, { identifier="setyujingdianyayuzhixiaxian", create_bin_func=control_setyujingdianyayuzhixiaxian, create_bin_args={nil,nil}, parser_bin_func=response_setyujingdianyayuzhixiaxian } } ------------------------------------- [串口 -> 平台] 解析输出 ------------------------------------- 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) sys.publish("MQTT_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) sys.publish("MQTT_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 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 ISO8601ToTime(datetime) log.info("ISO8601",datetime) local pattern = "(%d+)%-(%d+)%-(%d+)T(%d+):(%d+):(%d+)%.*" local xyear,xmonth,xday,xhour,xminute,xseconds,xmillies,xoffset = datetime:match(pattern) local convertedTimestamp = os.time({day = xday,month = xmonth,hour = xhour,year = xyear,sec = xseconds,min = xminute})-- return convertedTimestamp end local function AppHeartBeat(gw_sn,template_id) topic = string.format( "ipc/%s/script/%s/heartbeat",gw_sn,template_id) obj_send={} json_str = cjson.encode(obj_send) -- mqtt_client:publish(topic, json_str) sys.publish("MQTT_PUBLISH",topic, json_str) end local function ServiceCall(gw_sn,dev_sn,port,identifier,param) topic = string.format( "ipc/%s/%s/device/%s/data/Set_Rglt",gw_sn,port,dev_sn) obj_send={} obj_send['mi'] = mi -- obj_send['len'] = #data obj_send['identifier'] = identifier obj_send['sn'] = dev_sn obj_send['timestamp'] = os.time() if param then for k,v in pairs(param)do obj_send[k] = v end end json_str = cjson.encode(obj_send) log.info("ServiceCall",topic,json_str) -- mqtt_client:publish(topic, json_str) sys.publish("MQTT_PUBLISH",topic, json_str) end ---------------------------------------------------------- 外部消息---------------------------------------------------------- local function mqtt_session(id, addr, port, usr, pwd, subscribe, qos,on_mqtt_msg) sys.taskInit(function() local last = 0 while true do mqtt_client = MQTT.client.create(addr, port, on_mqtt_msg) local status = mqtt_client:connect("mandun") if not status then mqtt_client:subscribe(subscribe) while(true) do if os.difftime(os.time(),last) >= 5 then last = os.time() log.info("mqtt","check myself!!!!") AppHeartBeat(gw_sn,nodes[1].template_id) end sys.wait(100) local ret = mqtt_client:handler() if ret then log.info("mqtt","MQTT client was disconnect!") mqtt_client:destroy() break end end else sys.wait(5000) end 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_loop(nodes_cfg,sn) log.info("main_loop", "=============== MANDUN TEMPLATE 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_loop(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 subscribe={"ipc/+/+/device/+/data/raw_data","ipc/+/+/device/+/data/Set_Rglt"} -- subscribe[1] = string.format("ipc/+/+/device/+/raw_data") -- subscribe[1] = string.format("ipc/+/+/device/+/data/+") mqtt_session("internal","localhost", 1883 , nil, nil, subscribe, 0, on_mqtt_msg) local function conver2mosquito(gw_sn,gw_port,dev_sn,identifier,param) if not gw_sn or not gw_port or not dev_sn then return end local topic = string.format( "ipc/%s/%s/device/%s/data/Set_Rglt",gw_sn,gw_port,dev_sn) obj_send={} obj_send['mi'] = 0 obj_send['identifier'] = identifier obj_send['sn'] = dev_sn for k,v in pairs(param)do obj_send[k] = v end local payload = string.gsub(cjson.encode(obj_send),"\"","\\\"",99999) return topic,payload end -- schedule crontab sys.taskInit(function() while true do local result, dev_sn, data = sys.waitUntil("SCHEDULE_PERIOD", 1000) if result and dev_sn and data and data.items then --删除所有未执行调度器 log.info("schedule crontab","Schedule crontab") local cmd = "(crontab -l | grep -v \"mosquitto_pub\") | crontab -" -- log.info("Remove schedule crontab",cmd) os.execute(cmd) for _,item in ipairs(data.items) do if item and item.rule and item.identifier and item.msg then local topic,payload = conver2mosquito(gw_sn,FetchPortBySN(dev_sn),dev_sn,item.identifier,item.msg) cmd = string.format("cron_job='%s mosquitto_pub -t %s -m \"%s\"' ; ( crontab -l | grep -v \"$cron_job\"; echo \"$cron_job\" ) | crontab -",item.rule,topic,payload) os.execute(cmd) log.info("schedule crontab","Schedule crontab depoly Time ",time,cmd) end end end end end) -- schedule oneshot sys.taskInit(function() while true do local result, dev_sn ,data = sys.waitUntil("SCHEDULE_ONESHOT", 1000) if result and dev_sn and data and data.items then --删除所有未执行调度器 local time = os.date("%H:%M %Y-%m-%d",os.time()) log.info("schedule oneshot","Current time",time) local cmd = "atq | awk '{print $1}' | tr \"\\n\" \" \" | xargs atrm" -- log.info("Remove schedule",cmd) os.execute(cmd) for _,item in ipairs(data.items) do if item and item.date and item.identifier and item.msg then local topic,payload = conver2mosquito(gw_sn,FetchPortBySN(dev_sn),dev_sn,item.identifier,item.msg) -- HH:MM YYYY-MM-DD -- local time = os.date("%H:%M",ISO8601ToTime(item.date) ) local time = os.date("%H:%M %Y-%m-%d",ISO8601ToTime(item.date)) -- local time = "now" cmd = string.format("echo 'mosquitto_pub -t %s -m \"%s\"' | at %s",topic,payload,time) os.execute(cmd) log.info("schedule oneshot","Schedule atd depoly Time",time,cmd) end end end end end) -- mqtt 发送 sys.taskInit(function() while true do local ret,topic,payload = sys.waitUntil("MQTT_PUBLISH", 1000) if ret and topic then mqtt_client:publish(topic, payload) end end end) --启动系统框架 sys.init(0, 0) sys.run() end -- sn="218805000001" -- nodes='[{"connect_port":"RS485_2","depth":0,"product_key":"34851808","sn":"202007220001","template_id":"34851808","term_addr":"03"}]' -- main_loop(nodes,sn)