require "log" require "sys" require "utils" require "patch" require "pack" require "bit" local siganl = require 'posix.signal' local syswait = require 'posix.sys.wait' local unistd = require 'posix.unistd' local lpack = require("lua_pack") local internal_api = require("internal_api") LOG_LEVEL = log.LOGLEVEL_DEBUG --[[ 与PLC通讯状态==正常&&(急停开关==断开&&水浸开关==断开&&(门禁1开关==闭合&&门禁2开关==闭合)&&总开关==闭合) 急停 TAG_6AB 水浸 TAG_579 门禁 TAG_5DD 总开关 TAG_6B5 PLC 故障 TAG_4B5 BMS 电柜隔离开关故障 TAG_9DAE 消防 报警等级 TAG_4CA 烟雾 TAG_4C6 ]] --[[ # 发布日期:20221019 # 模块功能: - 防逆流电表通讯故障 - BMS禁充,禁放 - 交流配电总开关,闭合干接点 - 交流配电电表功率 > 5kw - UPS 工作状态开断干接点 - 消防故障开断干接点 - 电池舱空调通讯故障 - 电池舱空调故障反馈闭合干接点 - 弧光检测告警开断干接点 - 急停反馈闭合干接点 - 喷洒状态反馈开断干接点 - 消防火警闭合干接点 - 安全继电器动作信号干接点 # 应用项目: - 清安 ]] local sub_nodes = {{ gw = "21881EFF002B", sn_suffix = "21881EFF002B_PLC", grp_idx = nil, ident_list = {"Read_telemeter"}, offline_enable = true, }, { gw = "21881EFF002B", sn_suffix = "21881EFF002B_gridmeter", grp_idx = nil, -- nil:不分组总数据,1:第一组,2:第二组, ident_list = {"yougongzongdianliang", "sanxiangcv"}, -- 关心的采集服务, offline_enable = false, }, { gw = "21881EFF002B", sn_suffix = "21881EFF002B_Rack1_1", grp_idx = nil, -- nil:不分组总数据,1:第一组,2:第二组, ident_list = {"hardwarefault"}, -- 关心的采集服务 offline_enable = false, }, { gw = "21881EFF002B", sn_suffix = "21881EFF002B_aircond", grp_idx = nil, ident_list = {"readon"}, offline_enable = false, },{ gw = "21881EFF002B", sn_suffix = "21881EFF002B_DI", grp_idx = nil, -- nil:不分组总数据,1:第一组,2:第二组, ident_list = {"alarmdata"}, -- 关心的采集服务 offline_enable = true, }} local gw_sn local nodes_info -- SN:object local g_nodes -- 全局node_cfg local g_grp = { --[[{data={},ctrl={},status={},protect={},warn={}}]] } local g_total = { data = {}, ctrl = {}, status = {}, protect = {}, warn_type = {}, warn = {} } local function protect_level_get(g_total, g_grp) assert(g_total.protect) local protect_level = 0 if g_grp and g_grp.protect and g_grp.protect.level then if not g_total.protect.level then protect_level = g_grp.protect.level else protect_level = g_total.protect.level > g_grp.protect.level and g_total.protect.level or g_grp.protect.level end else protect_level = g_total.protect and g_total.protect.level and g_total.protect.level or 0 end return protect_level end local function init_grp_structure(grp_obj, sub_nodes) assert(grp_obj) assert(sub_nodes) local size = 0 for _, v in ipairs(sub_nodes) do if v.grp_idx then if v.grp_idx > size then size = v.grp_idx end end end if size == 0 then size = 1 end for grp_idx = 1, size do local data_space if not grp_obj[grp_idx] then grp_obj[grp_idx] = {} end data_space = grp_obj[grp_idx] if not data_space.protect then data_space.protect = {} end if not data_space.data then data_space.data = {} end if not data_space.ctrl then data_space.ctrl = {} end if not data_space.status then data_space.status = {} end if not data_space.warn_type then data_space.warn_type = {} end if not data_space.warn then data_space.warn = {} end end end -- 数据采集 local function pp_msg(topic, payload) local currTime = os.time() -- print("pp2app",topic,payload) local obj = cjson.decode(payload) local tags = obj.tags ---- 先对数据做汇聚 if not obj.sn then log.error(nil, "Invalid device sn") return end local node_info = nodes_info[obj.sn] if not node_info then return end local found for _, ident in ipairs(node_info.ident_list) do if obj.identifier == ident then found = true break end end if not found then --log.error(nil, string.format("Invalid identifier:%s in sn:%s", obj.identifier, obj.sn)) end local grp_idx = node_info.grp_idx local data_space if grp_idx then data_space = g_grp[grp_idx] else data_space = g_total end if data_space then -- 更新数据 for msg_field, msg_value in pairs(tags) do -- 遍历上来的变量 local data = data_space.data data[msg_field] = msg_value -- log.debug(nil,msg_field,msg_value) end -- 更新状态 node_info.recv_time = os.time() sys.publish("REAL_DATA_UPDATE", grp_idx) end -- log.debug(nil,cjson.encode(g_total)) end -- 控制命令 local function iot_msg(topic, payload) local currTime = os.time() -- print("pp2app",topic,payload) local obj = cjson.decode(payload) if obj.identifier == "protect_level_set" then local changed local last_level = g_total.protect.level if obj.protect_level then g_total.protect.level = obj.protect_level end if last_level ~= g_total.protect.level then changed = true end for grp_idx = 1, #g_grp do local data_space = g_grp[grp_idx] if not data_space then break end local key = string.format("grp%d_protect_level", grp_idx) local data = data_space.protect if obj[key] then data.level = obj[key] end if obj.grps_protect_level and obj.grps_protect_level[grp_idx] then local grp_protect = obj.grps_protect_level[grp_idx] local last_level = data.level if grp_protect.level then data.level = grp_protect.level end if last_level ~= data.level then changed = true end end end if changed then sys.publish("PROTECT_LEVEL_CHANGED") end -- log.debug(nil,g_total.protect.level,cjson.encode(g_grp[1])) end -- internal_api.pp2north(gw_sn,obj.sn,"NULL",obj.identifier,{},obj.mi) -- added by yut: renew topic and send to external broker local topic = string.format("ems/%s/service/set", obj.sn) local obj_send = {} obj_send['mi'] = obj.mi or 0 obj_send['identifier'] = obj.identifier obj_send['sn'] = obj.sn -- obj_send['port'] = "NULL" obj_send['time'] = os.time() -- obj_send['disable_report'] = 0 -- obj_send['report_period'] = 0 -- obj_send['data_type'] = 2 -- obj_send['userParam'] = "" -- obj_send['tag_node']=cjson.encode({}) -- table.insert(obj_send,data) json_str = cjson.encode(obj_send) mqtt_external_client:publish(topic, json_str) -- ended by yut end local function total_warn_level_report() while true do sys.waitUntil("WARN_LEVEL_CHANGED", 1000) local obj = { fault_level = g_total.warn.level and g_total.warn.level or 0, grps_fault_level = {} } for grp_idx, grp_obj in ipairs(g_grp) do table.insert(obj.grps_fault_level, { level = grp_obj.warn.level and grp_obj.warn.level or 0 }) end for i, v in ipairs(g_nodes) do -- internal_api.pp2north(gw_sn,v.sn,"NULL","fault_level_report",obj,0) local topic = string.format("/ems/%s/service", v.sn) local obj_send = {} obj_send['tags'] = obj obj_send['mi'] = obj.mi or 0 obj_send['identifier'] = "fault_level_report" obj_send['sn'] = v.sn obj_send['time'] = os.time() json_str = cjson.encode(obj_send) log.info("warn pack:", json_str) mqtt_external_client:publish(topic, json_str) end sys.wait(3000) end end local function total_warn_calculate() -- local grid_meter_sn = string.format("%s%s", gw_sn, "_gridmeter") -- if not nodes_info[grid_meter_sn] then -- nodes_info[grid_meter_sn] = {} -- end -- nodes_info[grid_meter_sn].recv_time = os.time() while true do local ret, grp_idx = sys.waitUntil("REAL_DATA_UPDATE", 2000) local now = os.time() if not grp_idx then -- 没有grp代表总数据,或者超时 local data_space = g_total local is_offline if data_space then local last_value = data_space.warn.level -- log.info(nil,"=========",nodes_info["21881EFF002B_DI"],cjson.encode(nodes_info)) for sn, v in pairs(nodes_info) do -- log.info(nil,string.format("dev[%s] diff = %d offline",sn,os.difftime(now, v.recv_time))) if not v.recv_time then v.recv_time = os.time() end if os.difftime(now, v.recv_time) >= 10 then -- if os.difftime(now, v.recv_time) >= 3 then v.offline = true else v.offline = false end -- log.info(nil,"=========",nodes_info["21881EFF002B_DI"],cjson.encode(nodes_info)) -- 状态改变通知其他模块上报 for sn, v in pairs(nodes_info) do if v.offline_enable and v.offline then is_offline = true break end end end -- log.info(nil,"=========",nodes_info["21881EFF002B_DI"],cjson.encode(nodes_info)) -- 急停 TAG_6AB --TAG_77B -- if data_space.data.TAG_EAA2 == nil then data_space.data.TAG_EAA2 = 0 end -- if data_space.data.TAG_EAA2 > 22 and data_space.data.TAG_EAA2 < 25 then data_space.warn.level = 3 end -- if data_space.data.TAG_EAA2 == 18 then data_space.warn.level = 3 end -- if data_space.data.TAG_EAA2 > 3 and data_space.data.TAG_EAA2 < 18 then data_space.warn.level = 2 end -- if data_space.data.TAG_EAA2 > 18 and data_space.data.TAG_EAA2 < 23 then data_space.warn.level = 2 end -- if data_space.data.TAG_EAA2 == 15 then data_space.warn.level = 2 end -- if data_space.data.TAG_EAA2 > 0 and data_space.data.TAG_EAA2 < 4 then data_space.warn.level = 1 end if data_space.data.TAG_77B == 1 or -- 急停反馈闭合干接点 DI 1915 TAG_77B data_space.data.TAG_776 == 1 or -- 消防喷洒状态反馈开断干接点 DI 1910 TAG_776 data_space.data.TAG_779 == 1 -- 消防火警闭合干接点 1913 TAG_779 then data_space.warn.level = 3 log.warn(nil,"level:",data_space.warn.level,cjson.encode(data_space.data)) -- 交流配电总开关位置反馈闭合干接点 1920 TAG_780 elseif data_space.data.TAG_780 == 1 or -- UPS正常工作状态开断干接点 1922 TAG_782 data_space.data.TAG_782 == 1 or -- 消防故障开断干接点 1911 TAG_777 data_space.data.TAG_777 == 1 or -- 弧光检测告警开断干接点 1916 TAG_77C data_space.data.TAG_77C == 1 or data_space.data.TAG_579 == 1 or -- 水浸 data_space.data.TAG_5DD == 1 or -- 门禁 data_space.data.TAG_6B5 == 1 or -- 总开关 data_space.data.TAG_9DAE == 1 or -- 电柜隔离开关故障 data_space.data.TAG_4CA == 1 or --消防报警等级 data_space.data.TAG_4C6 == 1 or --烟雾状态 data_space.data.TAG_EAA2 ~= 0 or is_offline then -- log.info("==========",is_offline) data_space.warn.level = 2 log.warn(nil,"level:",is_offline,data_space.warn.level,cjson.encode(data_space.data)) else data_space.warn.level = 0 end if last_value ~= data_space.warn.level then sys.publish("WARN_LEVEL_CHANGED") end end end end end local function grp_warn_calculate(idx) while true do local ret, grp_idx = sys.waitUntil("REAL_DATA_UPDATE", 2000) if ret and grp_idx and grp_idx == idx then local data_space = g_grp[grp_idx] if data_space then local last_value = data_space.warn.level if g_total.data[string.format("TAG_%3X", 0x77D + grp_idx)] == 1 -- DI模块放在总数据,不在组数据 -- 电堆n空调故障反馈开断干接点 1918 TAG_77E then data_space.warn.level = 2 else data_space.warn.level = 0 end -- 状态改变通知其他模块上报 if last_value ~= data_space.warn.level then sys.publish("WARN_LEVEL_CHANGED", idx) end end end end end local function total_protect_process(idx) local last_protect_level local last_warn_flag local action_step_idx = 0 local do_sn = string.format("%s%s", gw_sn, "_DO") local run_indicator local warn_indicator local fault_indicator while true do local now = os.time() local protect_level = g_total.protect.level and g_total.protect.level or 0 if last_protect_level ~= protect_level then action_step_idx = 0 last_protect_level = protect_level end if action_step_idx == 0 then if protect_level > 0 then log.warn(nil, string.format("System failure, level:%d", protect_level)) fault_indicator = 1 run_indicator = 0 else -- 没有保护,意味着没触发故障 log.warn(nil, string.format("System running")) run_indicator = 1 fault_indicator = 0 end -- UPS电池总压低告警开断干接点 DI 1921 TAG_781 -- 消防预警开断干接点 DI 1912 TAG_778 -- 弧光检测告警开断干接点 DI 1916 TAG_77C -- UPS旁路告警开断干接点 DI 1923 TAG_783 local warn_indicator = g_total.data.TAG_781 == 1 or g_total.data.TAG_778 == 1 or g_total.data.TAG_77C == 1 or g_total.data.TAG_783 == 1 if last_warn_flag ~= warn_indicator then last_warn_flag = warn_indicator if warn_indicator == 1 then log.warn(nil, string.format("System exited warning"), g_total.data.TAG_781, g_total.data.TAG_778, g_total.data.TAG_77C, g_total.data.TAG_783) end end end if run_indicator == 1 then -- 储能系统运行指示干接点 DO 1963 TAG_7AB internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AB"] = 1 }) -- 储能系统故障指示(声光报警)干接点 DO 1965 TAG_7AD internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AD"] = 0 }) end if warn_indicator == 1 then -- 储能系统报警指示干接点 DO 1964 TAG_7AC internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AC"] = 1 }) else -- 储能系统报警指示干接点 DO 1964 TAG_7AC internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AC"] = 0 }) end if fault_indicator == 1 then -- 储能系统运行指示干接点 DO 1963 TAG_7AB internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AB"] = 0 }) -- 储能系统故障指示(声光报警)干接点 DO 1965 TAG_7AD internal_api.call_dev_service(gw_sn, do_sn, "NULL", "failure", { ["TAG_7AD"] = 1 }) end sys.waitUntil("PROTECT_LEVEL_CHANGED", 30000) end end local function grp_protect_process(idx) local action_steps = {{ desc = "Breaker", cmp_obj = g_total.data, cmp_tag = "pcs_run_state", equal_v = 0, send_param = { sn = string.format("%s%s", gw_sn, "_DO"), -- TODO 需要等设备配置号了修改 identifier = "Grid_DO2", -- 动态计算TAG号 -- 并网点断路器1远程分闸干接点 DO 1960 TAG_7A8 -- 并网点断路器2远程分闸干接点 DO 1961 TAG_7A9 tag = string.format("TAG_%03X", 0X7A8 + idx - 1), value = 0 } }, { desc = "UPS", cmp_obj = g_total.data, cmp_tag = "bms_switch_on_state", equal_v = 0, send_param = { sn = string.format("%s%s", gw_sn, "_UPS"), -- TODO 需要等设备配置号了修改 identifier = "shutdown", -- 动态计算TAG号 -- UPS切载干接点 DO 1962 TAG_7AA -- 这里放在总的保护处理函数中也可以 tag = "TAG_7AA", value = 0 } }} local last_protect_level local action_step_idx = 0 while true do assert(idx) sys.waitUntil("PROTECT_LEVEL_CHANGED", 3000) -- 执行命令的序号 local protect_level = protect_level_get(g_total, g_grp[idx]) if last_protect_level ~= protect_level then action_step_idx = 0 last_protect_level = protect_level end if protect_level == 3 then if action_step_idx < #action_steps then local step = action_steps[action_step_idx + 1] assert(step.cmp_obj) assert(step.cmp_tag) assert(step.equal_v) assert(step.send_param) assert(step.send_param.sn) assert(step.send_param.identifier) assert(step.send_param.tag) assert(step.send_param.value) local last_time = os.time() local level_chenaged -- 测试代码 sys.timerStart(function () g_total.data.pcs_run_state = 0 end,3000) while os.difftime(os.time(), last_time) < 5000 do if step.cmp_obj[step.cmp_tag] == step.equal_v then break end if sys.waitUntil("PROTECT_LEVEL_CHANGED", 500) then level_chenaged = true break end end -- 等待过程中,保护等级变更 if not level_chenaged then internal_api.call_dev_service(gw_sn, step.send_param.sn, "NULL", step.send_param.identifier, { [step.send_param.tag] = step.send_param.value }) action_step_idx = action_step_idx + 1 log.info(nil, string.format("grp[%d] execute step[%d]:%s!", idx - 1, action_step_idx, step.desc)); end end end end end function update_nodes_info(sub_nodes) assert(sub_nodes) local nodes_info = {} for i, v in ipairs(sub_nodes) do if v.sn_suffix then nodes_info[v.sn_suffix] = v nodes_info[v.sn_suffix].recv_time = nil end end return nodes_info end local nodes_cfg = '[{"connect_port":"VINTF_TCP","depth":0,"product_key":"206627059","sn":"21881200000E_loadmetervir","template_id":"206627059"}]' function main(nodes_cfg, sn) log.info("main_loop", "=============== Guard System of EMS 1.0.0 ============") gw_sn = sn if not gw_sn then gw_sn = io.getGatewayID() end g_nodes = cjson.decode(nodes_cfg) for i, v in ipairs(g_nodes) do log.info("g_node _nodes_cfg,", v) end -- log.info("g_nodes:",g_nodes) init_grp_structure(g_grp, sub_nodes) nodes_info = update_nodes_info(sub_nodes) log.info("sub_nodes:", sub_nodes) log.info("nodes_info:", nodes_info) -- 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) -- added by yut local subscribe = internal_api.get_subscribes(gw_sn,sub_nodes,g_nodes) local subscribe = {} for i, v in ipairs(sub_nodes) do local topic = string.format("/ems/%s/service", v.sn_suffix) table.insert(subscribe, topic) log.info("sub topic", topic) end mqtt_external_client = internal_api.mqtt_session("115.238.104.202", 33004, "localuser", "dywl@galaxy", subscribe, function(topic, payload) --[[ if string.find(topic,"/data/service/") then pp_msg(topic,payload) --如果是原始数据一定是南向数据 elseif string.find(topic,"/data/Set_Rglt") then iot_msg(topic,payload) end ]] if string.find(topic, "/service/set") then iot_msg(topic, payload) -- 如果是原始数据一定是南向数据 elseif string.find(topic, "/service") then pp_msg(topic, payload) end end) sys.taskInit(total_warn_calculate) sys.taskInit(total_warn_level_report) sys.taskInit(total_protect_process) for grp_idx = 1, #g_grp do sys.taskInit(grp_warn_calculate, grp_idx) sys.taskInit(grp_protect_process, grp_idx) end -- 启动系统框架 sys.init(0, 0) sys.run() end main(arg[1], arg[2]) -- main(nodes_cfg, "2188120000E8")