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" RACK_PATH = "/app/" DB_FILE = RACK_PATH .. "racks.db" FTP_PATH = '/www/tf/' -- FTP_PATH = '' OSC_RECORD_ENABLE = false -- PACKAGE_TIME = 30 PACKAGE_TIME = 8*3600 -- DATA_TYPE_PROPERTY = 0, -- DATA_TYPE_EVENT, -- DATA_TYPE_SERVICE, LOG_LEVEL = log.LOGLEVEL_TRACE local rack_number = 1 local slave_number = 1 local vol_number = 13 local temp_number = 13 local vol_total = slave_number*vol_number local temp_total = slave_number*temp_number local last_report = nil local gw_sn = nil local cluster_sn = nil local gw_port = "RS485_1" local nodes_status = {} local service_info = {} local period_list = {} local message_cache = {} local pack_data_obj={} local cluster_obj={} local slave_vol_obj={} local slave_temp_obj={} local pack_sample_obj={} local last_report={} local last_cmd = {time = 0, func = 0, addr = 0, regaddr = 0} ------------------------------------------------------------------ ---------------------------- 业务触发配置表描述 ------------------------- ------------------------------------------------------------------ ----! identifier 服务标识,通过这个标识可以合并不同的采样包 ----! read_write 描述操作的读写 ----! class 采集数据类型:alarm/info,用于标记数据记录到不同的cache ----! regs.reginfo 寄存器名称(或名称的数组),寄存器地址,寄存器格式化符号,数值缩放比例 ----! regs.batch 批量采样个数,批量采样自累加值,格式化输出:字符串/数组 ----! regs.map 数值缩放原始值最小,原始值最大,输出值最小,输出值最大 ----! regs.level 报警等级(或者等级的数组) ------------------------------------------------------------------ global_point_table = { { sample_interval = 30, report_interval = 0, identifier = "ClusterInformationRequest", read_write = "R", regs = { {reginfo = {"StatusofConnectionProcess", 0x1001, ">S"}}, {reginfo = { { "EnableStatus_Rack1", "EnableStatus_Rack2", "EnableStatus_Rack3", "EnableStatus_Rack4", "EnableStatus_Rack5", "EnableStatus_Rack6", "EnableStatus_Rack7", "EnableStatus_Rack8", "EnableStatus_Rack9", "EnableStatus_Rack10", "EnableStatus_Rack11", "EnableStatus_Rack12", "EnableStatus_Rack13", "EnableStatus_Rack14", "EnableStatus_Rack15", "EnableStatus_Rack16" }, 0x1002, ">S" } }, {reginfo = {"RequestStepin", 0x1003, ">S"}}, {reginfo = {"RequestQuitout", 0x1004, ">S"}}, {reginfo = {"StatusAllow", 0x1005, ">S"}}, {reginfo = {"SystemOperationState", 0x1006, ">S"}}, {reginfo = {"SystemChargeState", 0x1007, ">S"}}, {reginfo = {"BAUAlarmState", 0x1008, ">S"}}, {reginfo = {"RackWarning", 0x1009, ">S"}}, {reginfo = {"RackAlarm", 0x100a, ">S"}}, {reginfo = {"RackCriticalAlarm", 0x100b, ">S"}}, {reginfo = {"SystemTotalVoltage", 0x100c, ">S",0.1}}, {reginfo = {"SystemTotalCurrent", 0x100d, ">s",0.1}}, {reginfo = {"SystemSoc", 0x100e, ">S",0.1}}, {reginfo = {"SystemSoh", 0x100f, ">S",0.1}}, {reginfo = {"SystemInsulation", 0x1010, ">S"}}, {reginfo = {"SystemEnableChargeEnergy", 0x1011, ">S"}}, {reginfo = {"SystemEnableDiscargeEnergy", 0x1012, ">S"}}, {reginfo = {"SystemMaxChargeCurrent", 0x1013, ">S",0.1}}, {reginfo = {"SystemMaxDischargeCurrent", 0x1014, ">S",0.1}}, {reginfo = {"RackCurrentDifference", 0x1015, ">s",0.1}}, {reginfo = {"RackVoltageDifference", 0x1016, ">S",0.1}}, {reginfo = {"SystemMaxVolCellRackID", 0x1017, ">S"}}, {reginfo = {"SystemMaxVolCellSlaveID", 0x1018, ">S"}}, {reginfo = {"SystemMaxVolCellID", 0x1019, ">S"}}, {reginfo = {"SystemMaxCellVoltage", 0x101a, ">S"}}, {reginfo = {"SystemMinVolCellRackID", 0x101b, ">S"}}, {reginfo = {"SystemMinVolCellSlaveID", 0x101c, ">S"}}, {reginfo = {"SystemMinVolCellID", 0x101d, ">S"}}, {reginfo = {"SystemMinCellVoltage", 0x101e, ">S"}}, {reginfo = {"Systemaveragevoltage", 0x101f, ">S"}}, {reginfo = {"SystemMaxTempCellRackID", 0x1020, ">S"}}, {reginfo = {"SystemMaxTempCellSlaveID", 0x1021, ">S"}}, {reginfo = {"SystemMaxTempCellID", 0x1022, ">S"}}, {reginfo = {"SystemMaxCellTemperature", 0x1023, ">s",0.1}}, {reginfo = {"SystemMinTempCellRackID", 0x1024, ">S"}}, {reginfo = {"SystemMinTempCellSlaveID", 0x1025, ">S"}}, {reginfo = {"SystemMinTempCellID", 0x1026, ">S"}}, {reginfo = {"SystemMinCellTemperature", 0x1027, ">s",0.1}}, {reginfo = {"Systemaveragetemperature", 0x1028, ">s",0.1}} } }, { sample_interval = 30, report_interval = 0, identifier = "ClusterAlarmRequest", read_write = "R", regs = { { reginfo = { { "BAUAlarmStateInternalComm", "BAUAlarmStatePCSComm", "BAUAlarmStatePCSControl", "BAUAlarmStateEMSComm", "BAUAlarmStateRackVolDiff", "BAUAlarmStateRackCurrentDiff", "BAUAlarmStateBAUAbort", "BAUAlarmStateAirconditionComm" }, 0x1008, ">S" } }, { reginfo = { { "RackVolHighWarm", "RackVolLowWarm", "CellVolHighWarm", "CellVolLowWarm", "DsgOverCurrWarn", "ChgOverCurrWarn", "DsgTempHighWarn", "DsgTempLowWarn", "ChgTempHighWarn", "ChgTempLowWarn", "InsulationLowWarn", "TerminalTempHighWarm", "HVBTempHighWarm" }, 0x1009, ">S" } }, { reginfo = { { "RackVolHighAlarm", "RackVolLowAlarm", "CellVolHighAlarm", "CellVolLowAlarm", "DsgOverCurrAlarm", "ChgOverCurrAlarm", "DsgTempHighAlarm", "DsgTempLowAlarm", "ChgTempHighAlarm", "ChgTempLowAlarm", "InsulationLowAlarm", "TerminalTempHighWarn", "HVBTempHighAlarm" }, 0x100a, ">S" } }, { reginfo = { { "RackVolHighCriticalAlarm", "RackVolLowCriticalAlarm", "CellVolHighCriticalAlarm", "CellVolLowCriticalAlarm", "DsgOverCurrCriticalAlarm", "ChgOverCurrCriticalAlarm", "DsgTempHighCriticalAlarm", "DsgTempLowCriticalAlarm", "ChgTempHighCriticalAlarm", "ChgTempLowCriticalAlarm", "InsulationLowCriticalAlarm", "TerminalTempHighCriticalAlarm", "HVBTempHighCriticalAlarm" }, 0x100b, ">S" } } } },{ sample_interval = 5, -- sample_interval = 3600*24, identifier = "RealtimeGet", read_write = "R", regs = { { reginfo = {"year",0x0003,">S"}, }, { reginfo = {"month",0x0004,">S"}, }, { reginfo = {"day",0x0005,">S"}, }, { reginfo = {"hour",0x0006,">S"}, }, { reginfo = {"minit",0x0007,">S"}, }, { reginfo = {"second",0x0008,">S"}, } } }, } 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 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, 1, 1)) == "A" then len = tonumber(string.sub(format, 2)) else log.warn("format lenght", "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 = {} log.info("read responce:",string.toHex(raw_data)) local _, addr, func, regaddr, regsize, skip_len if os.difftime(os.time(), last_cmd.time) > 10 then log.warn("ex_unpack", "recv timeout") return nil end if last_cmd.func and last_cmd.func == 0x03 or last_cmd.func == 0x04 then _, addr, func, bytes = bunpack(raw_data, "CCC") regaddr = last_cmd.regaddr regsize = bytes / 2 skip_len = 3 -- elseif last_cmd.func and last_cmd.func == 0x10 and last_cmd.func == func then -- _, addr, func, regaddr, regsize = bunpack(raw_data, "CC>S>S") -- skip_len = 6 else log.warn("ex_unpack", "modbus last_cmd canot find", last_cmd.func) return end if last_cmd.func ~= func then log.warn("ex_unpack", "modbus function code", string.format("send cmd:0x%02X recv cmd:0x%02X", last_cmd.func, func)) return nil end -- if last_cmd.addr ~= addr then -- log.warn("ex_unpack", "modbus slave addr", string.format("send addr:0x%02X recv addr:0x%02X", last_cmd.addr, addr)) -- return nil -- end local crc = dyutils.CRC16(raw_data, #raw_data - 2) local _, crc_pack = bunpack(string.sub(raw_data, #raw_data - 1), ">S") if (crc ~= crc_pack) then log.warn("ex_unpack", string.format("Data crc check failed cal:0x%04X packet:0x%04X", crc,crc_pack)) return nil end if #raw_data - skip_len <= 2 then -- log.warn("ex_unpack",#raw_data,skip_len) return nil end local data = string.sub(raw_data, skip_len + 1, -3) log.info("485cache_resopnse", addr, func, regaddr, regsize,string.toHex(data)) return true, addr, func, regaddr, regsize, data end local function convert2json(src_identifier,regaddr, regsize, raw_data) local service_index -- 遍历寄存器得到服务模型 for i, v in pairs(global_point_table) do if v.identifier == src_identifier then local vv = v.regs[1] if vv.reginfo[2] == regaddr then service_index = i break end end -- for ii, vv in pairs(v.regs) do -- if vv.reginfo[2] == regaddr then service_index = i end -- end end if not service_index then log.warn("port2pp", "not found index in point table") return end local obj = {} local levels = {} local identifier local skip_reg for reg = regaddr, regaddr + regsize - 1 do -- 跳转寄存器一旦被赋值,只有到了跳转到指定的寄存器才继续,否则直接退出 if not skip_reg or reg == skip_reg then -- 逐个寄存器查找效率地下 skip_reg = nil local v = global_point_table[service_index] local format = nil local key = nil local scale = nil local level = nil local batch = nil identifier = v.identifier class = v.class -- 找出寄存器名称和长度 for ii, vv in pairs(v.regs) do if vv.reginfo[2] == reg then format = vv.reginfo[3] key = vv.reginfo[1] scale = vv.reginfo[4] level = vv.level batch = vv.batch break end end if not format or not key then log.warn("port2pp", string.format("register info not found,regaddr:0x%04X", reg or -1)) else -- 计算出字段长度 local len = FormatLenght(format) -- 拆分转为对象 log.info("485cache_raw_data",string.toHex(string.sub(raw_data, 1, len)), format,len, string.format("regaddr:0x%04X", reg)) if type(key) == "table" then local _, value = bunpack(raw_data, format) if not value then log.warn("port2pp", "value cannot convert") return end raw_data = string.sub(raw_data, len + 1, -1) for pos, sub_key in pairs(key) do if not pos or pos <= 0 then log.warn("port2pp", "bunpack size error", format) return end if sub_key then obj[sub_key] = bit.band(bit.rshift(value, pos - 1), 1) end if sub_key and level then levels[sub_key] = level end end elseif type(key) == "string" then if not batch then local _, value = bunpack(raw_data, format) if not value then log.warn("port2pp","value cannot convert in batch process.",string.format("startaddr 0x%04X,currentaddr0x%04X, size %d,setp %d",reg, currreg, size, step)) return end raw_data = string.sub(raw_data, len + 1, -1) skip_reg = reg + len / 2 if scale then if type(value) ~= "number" then log.warn("port2pp","value was except as a number type.") return end obj[key] = value * scale else if type(value) == "string" then value = string.match(value, "%w+") end obj[key] = value end else local size = batch[1] or 1 local step = batch[2] or 1 obj[key] = {} for i = 1, size, step do local currreg = reg + (i - 1) -- log.info("port2pp", "=================",string.format("startaddr 0x%04X,currentaddr 0x%04X, size %d,setp %d",reg,currreg,size,step)) local _, value = bunpack(raw_data, format) if not value then log.warn("port2pp","value cannot convert in batch process.",string.format("startaddr 0x%04X,currentaddr0x%04X, size %d,setp %d",reg, currreg, size, step)) return end raw_data = string.sub(raw_data, len + 1, -1) skip_reg = currreg + 1 if scale then table.insert(obj[key], value * scale) else table.insert(obj[key], value) end end end end end end end return identifier, obj, class, levels end local function sync_system_time(obj_src) if not (obj_src.year and obj_src.month and obj_src.day and obj_src.hour and obj_src.minit and obj_src.second) then return end local date = os.date("%Y-%m-%d %H:%M:%S", os.time({year =obj_src.year, month = obj_src.month, day = obj_src.day, hour = obj_src.hour, min = obj_src.minit, sec = obj_src.second})) local cmd = "date -s " .. "\"" .. date .. "\"" log.info("==============BMS DATE", cmd) os.execute(cmd) end local function port2pp(input) local obj = cjson.decode(input) if not obj or not obj.src_identifier then return end if not obj.data_b64 then return end local raw_data = base64.decode(obj.data_b64) if not raw_data then return end local result, addr, func, regaddr, regsize, data = ex_unpack(raw_data) local curr = os.time() if result then if addr > 1 then if obj.src_identifier == 'Slave_Voltage' then sys.publish("VOL_RESPONSE",obj.sn,data,obj.mi) elseif obj.src_identifier == 'Slave_Temp' then sys.publish("TEMP_RESPONSE",obj.sn,data,obj.mi) end return end local dev_sn = FetchSnByAddr(addr) if not dev_sn then log.warn("port2pp", "SN not found") return end if bit.rshift(func, 7) ~= 0 then -- 差错寄存器 log.warn("port2pp", "received a fail response !") return elseif func == 0x10 then -- 写寄存器响应 log.info("port2pp", "write register response.", string.format("regaddr:0x%04X,regsize:%d", regaddr, regsize)) elseif func == 0x03 or func == 0x04 then -- 读寄存器响应 log.info("port2pp", "read register response.", string.format("regaddr:0x%04X,regsize:%d", regaddr, regsize),obj.src_identifier) local identifier, out, class, levels = convert2json(obj.src_identifier,regaddr, regsize, data) if not identifier or not out then log.warn("port2pp", "object convert failed!") return end if obj.src_identifier == "ClusterInformationRequest" then if not out.SystemTotalVoltage or not out.SystemTotalCurrent then return end cluster_data_obj = out pack_sample_obj[obj.sn] = obj.time if (obj.mi and obj.mi > 0) or not last_report[obj.src_identifier] or os.difftime(curr,last_report[obj.src_identifier]) > 300 then internal_api.service_response(gw_sn, obj.sn, "RS485_1", obj.src_identifier, out, obj.mi) last_report[obj.src_identifier] = curr end elseif obj.src_identifier == "ClusterAlarmRequest" then internal_api.service_response(gw_sn, obj.sn, "RS485_1", obj.src_identifier, out, obj.mi) elseif identifier == "RealtimeGet" then sync_system_time(out) end else log.info("port2pp", "function code invalid", func) end end end function read_from_485cache(gw_sn,port,dev_sn,identifier,addr,regaddr,regsize,mi) -- local topic = "to_Front/218800002001_TCP50001/transprent" local hex = ex_pack(addr,0x4,regaddr,regsize) log.info("read data:",string.toHex(hex)) internal_api.pp2south(gw_sn,dev_sn,port,identifier,hex,0,mi) end function write_to_485cache(gw_sn,port,dev_sn,identifier,addr,regaddr,regsize,data) -- local topic = "to_Front/218800002001_TCP50001/transprent" -- local payload = ex_pack(addr,0x10,regaddr,regsize,data) -- log.info("write cache:",topic,string.toHex(payload)) -- client:publish(topic, payload) end local function iot2pp(json_str, len) local err_code = 0 -- 取出入参-- log.info("iot2pp", "input:" .. json_str) local obj = cjson.decode(json_str) local identifier = obj.identifier local server_period = obj.server_period local dev_sn = obj.sn -- 编码-- if obj.identifier == "null" then log.warn("iot2pp", "Need device sn in payload") return 0, "", 0 end -- 不能直接解析,并发往端口,使用脚本的处理的内容 if obj.identifier == "RackListRequest" then local out = {Racks = {}} for k, v in pairs(pack_data_obj) do if v then table.insert(out.Racks, {sn = k, sample = pack_sample_obj[k], tags = v}) end end if #out.Racks < rack_number then return end internal_api.service_response(gw_sn, obj.sn, "RS485_1", obj.identifier, out, obj.mi) return else local addrs = FetchAddrBySN(dev_sn) if not addrs then log.warn("iot2pp", "Can not found term_addr") return nil end local addr = tonumber(addrs, 16) -- local addr = 2 -- 找服务 for i, v in pairs(global_point_table) do if v.identifier == obj.identifier then if not v.regs then log.warn("iot2pp", "regs not found in service") return end local function segment_process(output, rw, identifier, mi) if rw == "W" then if not output.reg_data then log.warn("iot2pp", "not find any binary data") return end log.debug("write reg data:",string.toHex(output.reg_data)) write_to_485cache(gw_sn, gw_port,dev_sn,identifier, addr,output.regaddr, output.regsize,output.reg_data) elseif rw == "R" then read_from_485cache(gw_sn, gw_port,dev_sn,identifier, addr, output.regaddr,output.regsize,obj.mi) end service_info.time = os.time() service_info.identifier = identifier service_info.mi = mi end local last_reg = nil local temp = {} for ii, vv in ipairs(v.regs) do -- 新的地址开头 if not temp.regaddr then temp.regaddr = vv.reginfo[2] temp.regsize = 0 temp.reg_data = '' end -- 判断地址是否连续 if (vv.reginfo[2] - temp.regaddr ~= temp.regsize) then segment_process(temp, v.read_write, obj.identifier, obj.mi) -- 缓存状态重新赋值 temp.regaddr = vv.reginfo[2] temp.regsize = 0 temp.reg_data = '' end local batch = vv.batch if not vv.batch then -- 写数据才需要从报文中获取字段value if v.read_write == "W" then -- 取出并叠加 local value = obj[vv.reginfo[1]] if not value then log.warn("iot2pp","can not find the key from message,key:",vv.reginfo[1]) return end temp.reg_data = temp.reg_data ..bpack(vv.reginfo[3], value) end local datasize = FormatLenght(vv.reginfo[3]) / 2 temp.regsize = temp.regsize + datasize else -- 批量下发需要循环 local size = vv.batch[1] or 1 local step = vv.batch[2] or 1 if v.read_write == "W" then local values = obj[vv.reginfo[1]] if not values then log.warn("iot2pp","can not find the key from message,key:",vv.reginfo[1]) return end end for i = 1, size, step do if v.read_write == "W" then -- 取出并叠加 if not values[i] then log.warn("iot2pp","index was invalid from message,key:",vv.reginfo[1], "index:", i) return end temp.reg_data = temp.reg_data ..bpack(vv.reginfo[3],values[i]) end temp.regsize = temp.regsize + 1 end end end -- 根据读写分类调用底层接口 segment_process(temp, v.read_write, obj.identifier, obj.mi) output = {} end end end return -1, nil, 0 end local function racks_cal(json_str, len) local err_code = 0 -- 取出入参-- local obj = cjson.decode(json_str) local tags = obj.tag_node and cjson.decode(obj.tag_node) or nil pack_data_obj[obj.sn] = tags pack_sample_obj[obj.sn] = obj.time end function on_external_msg(topic, payload) if string.find(topic, "data/Set_Rglt") then iot2pp(payload, #payload) elseif string.find(topic, "data/raw_data") then ----二进制数据流解析 port2pp(payload) elseif string.find(topic, "data/service/RackInformationRequest") then racks_cal(payload) else 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 call_dev_service(gw_sn,dev_sn,port,identifier,param,mi) local topic = string.format( "ipc/%s/%s/device/%s/data/Set_Rglt",gw_sn,port,dev_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['requester'] = 'local' obj_send['time'] = os.time() json_str = cjson.encode(obj_send) log.info("call_dev_service",topic,json_str) client:publish(topic, json_str) end function main(nodes_cfg, sn) log.info("main_loop", "=============== BMS PROTOCAL 1.2 ============") gw_sn = sn if not gw_sn then gw_sn = io.getGatewayID() -- gw_sn = "21012C00FFFF" 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 "")) gps_sn = nodes[1].sn cluster_sn = string.format("%s%02d", gw_sn, 1) local obj = cjson.decode(nodes_cfg) if obj and #obj > 0 then if obj.sn == cluster_sn and obj.ext_data and #obj.ext_data >2 then -- "ext_data": "{\"cellNum\":13,\"rackNum\":1}", local ext = string.gsub(obj.ext_data,"\\\"","\"",99999) local ext_obj = cjson.decode(ext) if ext_obj and ext_obj.cellNum then slave_number = 1 -- 默认1 vol_number = ext_obj.cellNum temp_number = ext_obj.cellNum end if ext_obj and ext_obj.rackNum then rack_number = ext_obj.rackNum end log.info("main_loop","Reload config",string.format("rack_number:%d vol_number:%d temp_number:%d",rack_number,vol_number,temp_number)) end end table.insert(nodes, {sn = cluster_sn, term_addr = 0x01}) siganl.signal(siganl.SIGKILL, handler) siganl.signal(siganl.SIGHUP, handler) local subscribe = { string.format("ipc/%s/RS485_1/device/+/data/raw_data", gw_sn), string.format("ipc/%s/RS485_1/device/+/data/service/RackInformationRequest",gw_sn), string.format("ipc/%s/RS485_1/device/%s/data/Set_Rglt", gw_sn,cluster_sn), string.format("ipc/%s/RS485_1/device/+/data/service/Slave_Voltage",gw_sn), string.format("ipc/%s/RS485_1/device/+/data/service/Slave_Temp",gw_sn) } internal_api.mqtt_session("localhost", 1883, nil, nil, subscribe,on_external_msg) -- 等待RTC就绪 while os.time() < 1536713845 do sys.wait(1000) call_dev_service(gw_sn,dev_sn,"RS485_1","RealtimeGet") end sys.taskInit(function () local last = 0 while(true) do sys.wait(5000) if os.difftime(os.time(),last) > 3600*24*10 then last = os.time() local dev_sn = string.format( "%s01",gw_sn) call_dev_service(gw_sn,dev_sn,"RS485_1","RealtimeGet") end end end) sys.taskInit(function () while(true) do sys.wait(5000) local dev_sn = string.format( "%s01",gw_sn) call_dev_service(gw_sn,dev_sn,"RS485_1","ClusterInformationRequest") end end) sys.taskInit(function () while(true) do sys.wait(27000) local dev_sn = string.format( "%s01",gw_sn) call_dev_service(gw_sn,dev_sn,"RS485_1","ClusterAlarmRequest") end end) sys.taskInit(function () while(true) do sys.wait(300000) local dev_sn = string.format( "%s01",gw_sn) call_dev_service(gw_sn,dev_sn,"RS485_1","RackListRequest") end end) local function update_data(data,request_off,request_len,rack_obj,scale) local offset_data = 1 local offset = request_off local size = request_len if not data then log.warn("unpack", "Invalid params in function update_data:",data,request_off,request_len,rack_obj) return end for index = offset + 1,offset + size do local value if data then local raw_data = string.sub(data,offset_data,offset_data+1) _, value = bunpack(raw_data, ">s") else value = 0 end offset_data = offset_data + 2 rack_obj[index] = value*scale end end if OSC_RECORD_ENABLE then for i = 1,rack_number do local racksn = string.format( "%s%02d",gw_sn,i+1) local addr = i+1 sys.taskInit(function (dev_sn,addr) local regsize = 100 if not slave_vol_obj[dev_sn] then slave_vol_obj[dev_sn] = {} end local rack_obj = slave_vol_obj[dev_sn] while(true) do sys.wait(10000) local offset = 0 local total = vol_total local size = 0 while(true) do sys.wait(5000) if total == 0 then break elseif total >= regsize then size = regsize total = total - size else ---- 等于直接退出 size = total total = total - size end read_from_485cache(gw_sn,"RS485_1",dev_sn,"Slave_Voltage",addr,0x2400 + offset,size,10000+offset) local last = os.time() local ret,data,mi while true do if os.difftime(os.time(),last) > 10 then log.warn("slave_voltage", "wait timeout!") break end ret,resp_sn,data,mi = sys.waitUntil("VOL_RESPONSE",1000) if ret and data and resp_sn and resp_sn == dev_sn and mi == 10000+offset then if size ~= #data/2 then log.warn("slave_voltage",string.format( "request len:%d received len:%d",size,#data/2)) data = nil else break end end end update_data(data,offset,size,rack_obj,1) offset = offset + size end end end,racksn,addr) sys.taskInit(function (dev_sn,addr) local regsize = 100 if not slave_temp_obj[dev_sn] then slave_temp_obj[dev_sn] = {} end local rack_obj = slave_temp_obj[dev_sn] while(true) do sys.wait(10000) local offset = 0 local total = temp_total local size = 0 while(true) do sys.wait(5000) if total == 0 then break elseif total >= regsize then size = regsize total = total - size else ---- 等于直接退出 size = total total = total - size end read_from_485cache(gw_sn,"RS485_1",dev_sn,"Slave_Temp",addr,0x2600 + offset,size,20000+offset) local last = os.time() local ret,data,mi while true do if os.difftime(os.time(),last) > 10 then log.warn("slave_temp", "wait timeout!") break end ret,resp_sn,data,mi = sys.waitUntil("TEMP_RESPONSE",1000) if ret and data and resp_sn and resp_sn == dev_sn and mi == 20000+offset then if size ~= #data/2 then log.warn("slave_temp",string.format( "request len:%d received len:%d",size,#data/2)) data = nil else break end end end update_data(data,offset,size,rack_obj,0.1) offset = offset + size end end end,racksn,addr) end local function writelineArray(parent,title,tags,date,rackid,perffix) if not tags then return end if not parent.title then local line_names = "" for _,name in ipairs(title) do line_names = line_names .. name .. "," end for i = 1,#tags do line_names = line_names .. (perffix or "") .. i .. "," end parent.title = true parent.file.write(parent.file,line_names .. "\n") end local line_values = "" for _,name in ipairs(title) do if name == "日期时间" or name == "Date" then line_values = line_values .. date .. "," elseif name == "RackId" then line_values = line_values .. rackid .. "," end end for i = 1,#tags do line_values = line_values .. (tags[i] and tags[i] or "--") .. "," end parent.file.write(parent.file,line_values .. "\n") end sys.taskInit(function () sys.wait(1000) while true do local curr = os.time() local last_time = curr local last_save = curr local start_date = os.date("%Y-%m-%d_%H-%M-%S") local cluster_info = {title=nil, file = io.open( string.format( FTP_PATH.."Bank_%s.csv",start_date), "a+" )} local rack_info = {title=nil, file = io.open( string.format( FTP_PATH.."Rack_%s.csv",start_date), "a+" )} local body_vol = {title=nil, file = io.open( string.format( FTP_PATH.."CellVoltage_%s.csv",start_date), "a+" )} local body_temp = {title=nil, file = io.open( string.format( FTP_PATH.."CellTemp_%s.csv",start_date), "a+" )} while true do sys.wait(500) curr = os.time() if os.difftime(curr,last_save) > 5 then last_save = curr if os.difftime(curr,last_time) > PACKAGE_TIME then break end local date = os.date("%Y-%m-%d %H:%M:%S") ----- 保存总控概要 if cluster_data_obj then writeline(cluster_info,title_ClusterInformationRequest,cluster_data_obj,date,rackid,alias_ClusterInformationRequest) end for i = 1,rack_number do local racksn = string.format( "%s%02d",gw_sn,i+1) local rackid = i ----- 保存主控概要 if pack_data_obj[racksn] then writeline(rack_info,title_RackListRequest,pack_data_obj[racksn],date,rackid,alias_RackListRequest) end ----- 保存主控单体 if slave_vol_obj[racksn] and #slave_vol_obj[racksn] == (slave_number * vol_number) then writelineArray(body_vol,title_Body,slave_vol_obj[racksn],date,rackid,"#(mV)") end ----- 保存主控温度 if slave_temp_obj[racksn] and #slave_temp_obj[racksn] == (slave_number * temp_number) then writelineArray(body_temp,title_Body,slave_temp_obj[racksn],date,rackid,"#(°C)") end end end end cluster_info.file.close(cluster_info.file) rack_info.file.close(rack_info.file) body_vol.file.close(body_vol.file) body_temp.file.close(body_temp.file) end end) end -- 启动系统框架 sys.init(0, 0) sys.run() end -- main(arg[1], arg[2]) main("[{sn:123123:term_addr:121}]", "12313212")