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") lpack=require("lua_pack") require "internal_api" require"lsqlite3" require"nodes_sta" LOG_LEVEL = log.LOGLEVEL_TRACE local last_report = nil local gw_sn=nil local service_info = {} local period_list = {} local gw_port = "NULL" local last_cmd = {time=0,func=0,addr=0,regaddr=0} local function add_period_service(node_sn) for _,service in ipairs(global_point_table) do if service.read_write == "R" then if not service.sample_interval or service.sample_interval == 0 then break end local period={ identifier = service.identifier, sn = node_sn, intv_sample = service.sample_interval, -- 直接使用上报周期做采样 intv_report = service.report_interval, -- 直接使用上报周期做采样 last_sample = 0, last_report = 0, heartbeat = os.time() } table.insert(period_list,period) end end end local function del_period_service(node_sn) for k,v in ipairs(period_list)do if v and v.sn and v.sn == node_sn then table.remove(period_list,k,1) end end end local function UpdateAddrSN(sn,addr) local update_sn for _,v in pairs(nodes) do if v.sn == sn then v.term_addr = addr for k,vv in ipairs(period_list)do if vv and vv.sn and vv.sn == sn then vv.heartbeat = os.time() break end end for _,vv in pairs(nodes) do if vv.sn ~= sn and vv.term_addr == addr then vv.term_addr = nil end end return end end table.insert(nodes,{sn = sn,term_addr = addr}) add_period_service(sn) 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 FetchSnByAddr(addr) for _,v in pairs(nodes) do if v.term_addr == addr then return v.sn end end end local function DeleteAddrBySN(sn) for i,v in ipairs(nodes) do if v.sn == sn then del_period_service(sn) return table.remove( nodes, i, 1) end end end local function FormatLenght(format) local len if string.upper(format) == "C" then len = 1 elseif string.upper(string.sub(format,2,2)) == "S" then len = 2 elseif string.upper(string.sub(format,2,2)) == "I" then len = 4 elseif string.upper(string.sub(format,2,2)) == "D" then len = 8 elseif string.upper(string.sub(format,1,1)) == "A" then len = tonumber(string.sub(format,2)) else log.warn("port2pp", "format not support",format) return 0 end return len end --! @brief 模块功能:485cache 组包 --! @param addr: 设备地址 --! @param func: 功能码 --! @param regaddr: 寄存器地址 --! @param regsize: 寄存个数 --! @param data: 写操作时:寄存器个数 *2 字节;读操作时:0 --! @return --! @retval raw_data: 返回二进制数据 function ex_pack(addr,func,regaddr,regsize,data) log.info("485cache_api",addr,func,string.format( "0x%04X",regaddr) ,regsize,data and string.toHex(data)) local raw_data=nil if not addr or not func or not regaddr or not regsize then log.warn("ex_pack","修改设备层级") return nil end raw_data = bpack("CC>S>S",addr,func,regaddr,regsize) -- 有数据增加二进制数据 if data and #data > 0 then if func == 0x01 and #data ~= regsize/2 then log.warn("ex_pack",string.format("buff size and regsize are conflict, size:%d,regsize:%d",#data,regsize)) return nil end 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)) last_cmd = {time=os.time(),func=func,addr=addr,regaddr=regaddr} return raw_data end --! @brief 模块功能:485cache 拆包 --! @param raw_data: 输入二进制数据 --! @return --! @retval result: 返回结果 --! @retval level: 设备层级 --! @retval addr: 设备地址 --! @retval func: 功能码 --! @retval regaddr: 寄存器地址 --! @retval regsize: 寄存个数 --! @retval data: 写操作时:寄存器个数 *2 字节;读操作时:0 function ex_unpack(raw_data) --取出入参-- local obj={} local _,head,mac,type_d,fun,len = bunpack(raw_data,"C6CCC>S") local addr = string.sub(raw_data,12,15) if head ~= 0x68 then log.warn("ex_unpack","head code",string.format("cmd:0x%02X",head)) return nil end local _,stop = bunpack(string.sub(raw_data,#raw_data),"C") if stop ~= 0x16 then log.warn("ex_unpack","stop code",string.format("cmd:0x%02X",stop)) return nil end local _,sum_rcv = bunpack(string.sub(raw_data,#raw_data - 2,#raw_data - 1),">S") local sum = 0 for i=1,#raw_data -3 do local _,dat = bunpack(string.sub(raw_data,i,i),"C") sum = sum + dat end sum = bit.band(sum,0xffff) if (sum ~= sum_rcv) then log.warn("ex_unpack",string.format("Data crc check failed cal:0x%04X packet:0x%04X",sum,sum_rcv)) return nil end local data = string.sub(raw_data,12 + 4,-4) -- log.info("out response:",string.toHex(addr," "),string.toHex(data," ")) return true,string.toHex(addr),data end --! @brief 模块功能:内部服务上报,或者服务请求的应答,用于IOT侧上行数据 --! @param gw_sn: 网关SN --! @param dev_sn: 表示设备SN --! @param port: 设备port --! @param identifier: 服务标识 --! @param param: 数据的具体内容,类型是table --! @param mi: 消息流水号,默认为0 --! @param time: 采集时间,默认是当前时间 function service_response(gw_sn,dev_sn,port,identifier,param,mi,time) local topic = string.format( "/sys/WqHEsTmE/%s/service",gw_sn) 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"] = 2 obj_send["tags"] = param and type(param) == "table" and param or cjson.decode(param) json_str = cjson.encode(obj_send) log.info("service_response",topic,json_str) client:publish(topic, json_str) end local function conv2obj(sn,raw_data) local data local tab = { {key="batt",format = ">S",scale=0.1}, {key="coil1",format = ">S",scale=0.1}, {key=nil,format = ">S"}, {key="coil2",format = ">S",scale=0.1}, {key=nil,format = ">S"}, {key="coil3",format = ">S",scale=0.1}, {key=nil,format = ">S"}, {key="coil4",format = ">S",scale=0.1}, {key=nil,format = ">S"}, {key=nil,format = ">d"}, {key=nil,format = "A8"}, {key=nil,format = "A6"}, {key="devid",format = "A20"}, {key="active",format = "C"}, {key="interv",format = ">S"}, {key="ver",format = "C",scale=0.1}, {key="csq",format = "C"}, } local obj ={} for _,o in ipairs(tab) do local len = FormatLenght(o.format) local _,value = bunpack(raw_data,o.format) if o.key then if o.scale then obj[o.key] = value * o.scale else obj[o.key] = value end end raw_data = string.sub(raw_data,len + 1) end return obj end local function port2pp(socket,input) log.info("response:",string.format("socket:%s",socket),string.toHex(input)) local result,sn,raw_data = ex_unpack(input) if result then log.info("port2pp","read register response.",string.format("regaddr:%s",sn)) local obj = conv2obj(sn,raw_data) log.info("out",cjson.encode(obj)) service_response(gw_sn,sn,gw_port,"post",obj) nodes_sta.rx_add(sn) end end local function iot2pp(json_str, len) end function on_external_msg(topic, payload) if string.find(topic,"transparent") then ----二进制数据流解析 local obj = cjson.decode(payload) if obj and obj.payload and obj.socket then local hex = base64.decode(obj.payload) if not obj.socket or not obj.payload then log.warn("mqtt.recv","frame error") return end print("topic:",topic,"payload:",payload) port2pp(obj.socket,hex) end else log.info("due message:",topic, payload) end end local function handler() local pid, status, code = syswait.wait(-1, syswait.WNOHANG) log.error('system exit', pid, status, code) --打印不出来 os.exit(0) end function main(nodes_cfg,sn) log.info("main_loop", "=============== shuzhu-shangkanyuan RUN 1.1 ============") gw_sn = sn if not gw_sn then gw_sn = io.getGatewayID() end log.info("main_loop", string.format("Try Run At Shell\nsn='%s'\nnodes='%s'\nmain_loop(nodes,sn)",gw_sn or "NOT SN",nodes_cfg or "")) siganl.signal(siganl.SIGKILL,handler) siganl.signal(siganl.SIGHUP,handler) local subscribe={string.format("/sys/WqHEsTmE/%s/service/set",gw_sn),"F/shuzhushangkanyuan/+/transparent"} internal_api.mqtt_session("mqtt.lnxall.com",3883,"localuser","dywl@galaxy",subscribe,on_external_msg) sys.taskInit(function () while true do sys.wait(3000) local report,change = nodes_sta.info(60*60*3) if not last_report then last_report = 0 end if change or os.difftime(os.time(),last_report) >= 3600 then last_report = os.time() local str = cjson.encode(report) report = nil if str then local topic = string.format( "/sys/WqHEsTmE/%s/nodes_status",gw_sn) client:publish(topic, str) end end end end) --启动系统框架 sys.init(0, 0) sys.run() end main(arg[1],arg[2])