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 mqtt=require "mosquitto" local base64=require("base64") local dyutils=require("dyutils") local lpack=require("lua_pack") require "internal_api" require"lsqlite3" require"shmem" package.path = "/app/hmi4rack/script/?.lua;" local reg_config=require"reg_config" local nodes = reg_config and reg_config.nodes and reg_config.nodes or {} RACK_PATH = "/app/hmi4rack/" DB_FILE = RACK_PATH .. "database/".. "racks.db" CFG_FILE = RACK_PATH .. "config/".. "project_para.json" -- DATA_TYPE_PROPERTY = 0, -- DATA_TYPE_EVENT, -- DATA_TYPE_SERVICE, -- LOG_LEVEL = log.LOGLEVEL_DEBUG LOG_LEVEL = log.LOGLEVEL_INFO TIMEZOEN = 8*3600 VAILD_TIME = 200 DEF_RACK_NUM=18 DEF_PACK_NUM=12 DEF_BATT_NUM=240 DEF_TEMP_NUM=120 DEF_BATTPOLE_NUM=24 REAL_PASER_ENABLE = false OPTION_WRITE = "W" OPTION_READ = "R" -- 全局寄存器,需要表达式中引用 regs_space={} --寄存器地址<->寄存器当前的值 regs_time={} --寄存器地址<->寄存器更新时间 func_expr={} --表达式标识<->loadstring后的function enumstr_expr={} -- 表达式标识枚举转换<->loadstring后的function local gw_sn=nil local service_info = {} local period_list = {} local gw_port = "CANBUS_0" local pack_alarm_obj={} local pack_data_obj={} local pack_connect_obj={} local to_dev_mq={} -- 发送给设备的MQ local identifier_index={} -- 快速查表定位寄存器所在identifier index local sync_block_index={} -- 快速查表定位寄存器所在sync_block index local system_parameter = {} -- 系统管理参数 local racknum local date_synced = nil -------------- Short name mapping------------- format = string.format band = bit.band rshift = bit.rshift local function FetchScanRegBySN(sn) for _,v in pairs(nodes) do if v.sn == sn then return v.scan_regs end end end local function FetchDevPointBySN(sn) for _,v in pairs(nodes) do if v.sn == sn then return v.dev_point end end end local function FetchAddrBySN(sn) for _,v in pairs(nodes) do if v.sn == sn then return v.term_addr,v.dev_lev 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 spilt_addr_str(str) local raw_data = str local out while true do local seg_in= string.match(raw_data,"&%w+") if not seg_in then break end local reg=tonumber(string.sub(seg_in,2),16) raw_data = string.gsub(raw_data,seg_in,'') if not out then out={} end table.insert(out,reg) end return out end local function identifier_index_init() for _,node in ipairs(nodes) do local addrs,dev_lev = FetchAddrBySN(node.sn) if not addrs then log.warn("parser_sync_cache_table","Can not found term_addr") return nil end local addr = tonumber(addrs,10) if addr and addr > racknum then log.info(nil,string.format("addr %d and behind it be discard",addr)) break end local dev_sn = node.sn if not identifier_index[dev_sn] then identifier_index[dev_sn] = {} end local dev_identifier_index = identifier_index[dev_sn] local point_table = FetchDevPointBySN(node.sn) if not point_table then log.warn(nil, "not found point table with sn:",dev_sn) return end for i,v in ipairs(point_table) do for ii,vv in ipairs(v.regs) do local regaddr = vv.reginfo[2] local extrainfo = vv.extrainfo local batch = vv.batch if type(regaddr) == "number" then if not batch then dev_identifier_index[regaddr]={i,ii} -- log.debug("init",string.format("Initalized point table,node:%s reg 0x%04X:%d identifier:%s index:%s",dev_sn,regaddr,ii,v.identifier,i) ) else local size = batch[1] or 1 local step = batch[2] or 1 for offset = 1,size,step do local currreg = regaddr + (offset - 1) dev_identifier_index[currreg]={i,ii} -- log.debug("init",string.format("Initalized point table,node:%s reg 0x%04X:%d identifier:%s index:%s",dev_sn,regaddr,ii,v.identifier,i) ) end end elseif type(regaddr) == "string" then -- 解析字符串中的寄存器 local regs = spilt_addr_str(regaddr) for _,regaddr in ipairs(regs) do if regaddr then dev_identifier_index[regaddr]={i,ii} -- log.debug("init",string.format("Initalized point table,node:%s reg 0x%04X:%d identifier:%s index:%s",dev_sn,regaddr,ii,v.identifier,i) ) end end end if extrainfo then for _,obj in ipairs(extrainfo) do if obj[2] and type(obj[2]) == "string" then -- 解析字符串中的寄存器 local regs = spilt_addr_str(obj[2]) for _,regaddr in ipairs(regs) do if regaddr then dev_identifier_index[regaddr]={i,ii} -- log.debug("init",string.format("Initalized extrainfo table,node:%s reg 0x%04X:%d identifier:%s index:%s",dev_sn,regaddr,ii,v.identifier,i) ) end end end end end end 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("port2pp", "format not support",format) return 0 end return len end local function AlarmContentSort(a, b) -- 不要用 <= 会出bug return a.settime < b.settime end local function replyImmediateAlarm(gw_sn,dev_sn,port,identifier,param,mi) local obj={alarmlist={}} local sql sql = string.format("SELECT sn,rackid,level,settime,clrtime,alarmid,extra FROM alarm where clrtime is null group by rackid,alarmid order by id desc limit 0,200") log.info("db fetch immedia sql:",sql) for dev_sn,rackid,level,st,ct,alarmid,extra in db:urows(sql) do if not rackid or not st or not alarmid then log.error("db",db:errcode(),db:errmsg()) break end table.insert(obj.alarmlist,{sn=dev_sn,rackid=rackid,settime=st,alarmitem=alarmid,level=level,extra=extra}) end table.sort(obj,AlarmContentSort) internal_api.service_response(gw_sn,dev_sn,port,identifier,obj,mi) end local function replyHistoryAlarm(gw_sn,dev_sn,port,identifier,param,mi,page,size) if not page then page = 1 end if not size then size = 100 end local obj={alarmlist={},pages=1,total=0,size=size} local sql local total_item sql = string.format("SELECT count(*) FROM (SELECT * from alarm where clrtime > 0)") log.info("db fetch history sql:",sql) for count in db:urows(sql) do total_item = count break end if page <=0 then log.warn("page start from 1") return end if size <=0 then log.warn("size need than 0") return end obj.pages = math.ceil(total_item/size) obj.total = total_item local start,stop = (page-1)*size,page*size sql = string.format("SELECT sn,rackid,level,settime,clrtime,alarmid,extra FROM (SELECT * from alarm where clrtime > 0 order by id desc)a limit %d,%d",start,stop) log.info("db fetch history sql:",sql) for dev_sn,rackid,level,st,ct,alarmid,extra in db:urows(sql) do if not rackid or not st or not ct or not alarmid then log.error("db",db:errcode(),db:errmsg()) break end table.insert(obj.alarmlist,{sn=dev_sn,rackid=rackid,settime=st,clrtime=ct,alarmitem=alarmid,level=level,extra=extra}) end table.sort(obj,AlarmContentSort) internal_api.service_response(gw_sn,dev_sn,port,identifier,obj,mi) end local function replyAlarmCount(gw_sn,dev_sn,port,identifier,param,mi) local obj={alarmlist={}} local sql local real_alarm,hist_alarm sql = string.format("SELECT count(*) FROM alarm where clrtime is null") log.info("db fetch immedia sql:",sql) for count in db:urows(sql) do real_alarm = count break end sql = string.format("SELECT count(*) FROM (SELECT * from alarm where clrtime > 0)") log.info("db fetch history sql:",sql) for count in db:urows(sql) do hist_alarm = count break end internal_api.service_response(gw_sn,dev_sn,port,identifier,{real_alarm=real_alarm,hist_alarm=hist_alarm},mi) end local function replyCleanAlarm() log.warn("clean alarm database") db:exec('delete from alarm;') end local function replyCellInfoRequest(gw_sn,dev_sn,port,identifier,param,mi) internal_api.service_response(gw_sn,dev_sn,port,identifier,param,mi,true) end local function db_autoclean() for page_count in db:urows('PRAGMA page_count') do for page_size in db:urows('PRAGMA page_size') do local database_size = page_size * page_count log.info("database size:",database_size) if database_size > 50*1024*1024 then log.warn("database autoclear database") db:exec('delete from realdata where id in (select id from realdata order by id asc limit 1000 offset 0) ;') end end end -- db.exec('delete from realdata where id in (select id from realdata order by id asc limit 1000 offset 0) ;') end local function db_alarm_status_insert(stmt,dev_sn,rack_id,index,errormsg,level,extra) if not dev_sn or not rack_id or not errormsg then log.error("alarm","param failed",dev_sn,rack_id,level,errormsg,extra) return end log.debug(nil,"Rack",rack_id,"insert alarm:",errormsg,extra) stmt:bind_values(dev_sn,rack_id,os.time(),nil,0,level,errormsg,extra or "") stmt:step() if db:errcode() == sqlite3.ERROR then log.error("db_error",db:errmsg()) return end stmt:reset() end local function db_alarm_status_update(stmt,dev_sn,rack_id,index,errormsg) --更新数据库结束时间 if not dev_sn or not rack_id then log.error("alarm","param failed",dev_sn,rack_id,errormsg) return end log.debug(nil,"Rack",rack_id,"clear alarm:",errormsg) stmt:bind_values(os.time(),rack_id,errormsg) stmt:step() if db:errcode() == sqlite3.ERROR then log.error("db_error",db:errmsg()) return end stmt:reset() end local function db_init(path) db = sqlite3.open(path) if not db then log.warn("db_init","db open fail,file:",path) return end db:busy_timeout(5000) db:exec('select * from alarm where extra=""') if db:errcode() > 0 then log.warn("db", "Table of alram to remove as no extra",db:errcode()) db:exec('DROP TABLE alarm') end db:exec('CREATE TABLE alarm(id INTEGER PRIMARY KEY,sn VARCHAR(32),rackid INTEGER,settime INTEGER,clrtime INTEGER,count INTEGER,level INTEGER,alarmid VARCHAR(32),extra VARCHAR(32) NOT NULL)') for page_count in db:urows('PRAGMA page_count') do for page_size in db:urows('PRAGMA page_size') do local database_size = page_size * page_count log.info("database size:",database_size) if database_size > 50*1024*1024 then log.warn("database autoclear database") end end end stmt_insert_alarm = db:prepare('INSERT INTO alarm VALUES (NULL,?,?,?,?,?,?,?,?)') -- log.fatal("db_init",db:errmsg()) stmt_update_alarm = db:prepare('update alarm SET clrtime = ? WHERE id = (select MAX(id) AS max_id from alarm where rackid = ? and alarmid = ? and clrtime IS NULL)') -- db:exec('CREATE TABLE realdata(id INTEGER PRIMARY KEY,sn VARCHAR(32),rackid INTEGER,time INTEGER,identifier VARCHAR(64),report INTEGER,msg TEXT)') -- stmt_insert_data = db:prepare('INSERT INTO realdata VALUES (NULL,?,?,?,?,?,?)') local sql = string.format("SELECT rackid,alarmid FROM alarm where clrtime is null group by rackid,alarmid order by id desc limit 0,200") for rack_id,alarm_key in db:urows(sql) do if rack_id and alarm_key then if not pack_alarm_obj[rack_id] then pack_alarm_obj[rack_id] = {} end table.insert(pack_alarm_obj[rack_id],alarm_key) end end -- local tmp = pack_alarm_obj[rack_id][1] end --! @brief 模块功能:485cache 组包 --! @param level: 设备层级 --! @param addr: 设备地址 --! @param func: 功能码 --! @param regaddr: 寄存器地址 --! @param regsize: 寄存个数 --! @param data: 写操作时:寄存器个数 *2 字节;读操作时:0 --! @return --! @retval raw_data: 返回二进制数据 local function ex_pack(level,addr,func,regaddr,regsize,data) -- log.debug("485cache_api",level,addr,func,string.format( "0x%04X",regaddr) ,regsize,data and string.toHex(data)) local raw_data=nil if not level or not addr or not func or not regaddr or not regsize then log.warn("ex_pack",level,addr,func,regaddr,regsize,data) return nil end raw_data = bpack("CCC>S>S",level,addr,func,regaddr,regsize) -- 有数据增加二进制数据 if data and #data > 0 then if func == 0x01 and #data/2 ~= regsize 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)) 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={} _,level,addr,func,regaddr,regsize = bunpack(raw_data,"CCC>S>S") 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("ex_unpack",string.format("Data crc check failed cal:0x%04X packet:0x%04X",crc,crc_pack)) return nil end if #raw_data <= 9 then log.info("ex_unpack",string.format("Frame length error ,raw_data length: %d",#raw_data)) return nil end local data = string.sub(raw_data,8,-3) log.debug("device_resopnse",level,addr,func,regaddr,regsize,string.toHex(data)) return true,level,addr,func,regaddr,regsize,data end local function expr2function(term_addr,expr,index,key) local raw_data = index ---- 将配置中的&0x0011取地址符号,转换成lua查表方式space[0x0011] while true do local seg_in= string.match(raw_data,"&%w+") if not seg_in then break end local seg_out=string.format("(shmem.shmsingle(%d,%s))",term_addr,string.sub(seg_in,2)) raw_data = string.gsub(raw_data,seg_in,seg_out) end local exp = string.format("return %s",raw_data) local f = loadstring(exp) assert(f) local result,data = pcall(f) if result then expr[index] = f log.info(nil, string.format("Do expr. done as calling [%s] for key:%s",exp,key) ) else log.debug(nil, string.format("Do expr. failed as calling [%s] for key:%s",exp,key) ) log.warn(nil, data) end end local function enum2expr(expr,enumstr,key,value) local exp if type(value) == "number" then -- 只有数值才会比对枚举 exp = string.format("local tab=%s local find=0 for key,enum_v in pairs(tab) do if enum_v == %d then return key end end if find == 0 then return '--' end",enumstr,value) else exp = string.format("return '--'",enumstr,value) end local f = loadstring(exp) assert(f) local result,data = pcall(f) if result then expr[enumstr] = f log.debug(nil, string.format("Do expr. done as calling [%s] for key:%s",exp,key) ) else log.debug(nil, string.format("Do expr. failed as calling [%s] for key:%s",exp,key) ) log.warn(nil, data) end return result,data end local function convert_local_time_to_GMT(obj_src) for _,value in ipairs({"year","month","day","hour","minit","second"}) do if not obj_src[value] then log.warn(nil,"Conver fail",cjson.encode(obj_src)) return end if type(obj_src[value]) == "string" then log.warn(nil,"Cache invalid") return end end if (obj_src.year < 2000 or obj_src.month < 1 or obj_src.day < 1) then log.warn(nil,"Date format invalid") return end local local_ts_tb = {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 local_ts = os.time(local_ts_tb) if not local_ts then log.warn(nil,"conver fail",cjson.encode(local_ts_tb)) return end local gmt_ts = local_ts - TIMEZOEN local gmt_ts = os.date("%Y-%m-%d %H:%M:%S", gmt_ts) log.info("Local time : " .. os.date("%Y-%m-%d %H:%M:%S", local_ts)) log.info("GMT time : " .. gmt_ts) return gmt_ts 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 gmt_ts = convert_local_time_to_GMT(obj_src) if gmt_ts then local cmd = "date -s " .. "\"" .. gmt_ts .. "\"" os.execute(cmd) date_synced = true end end local function findIdentifierByregaddr(dev_sn,regaddr) local point_table = FetchDevPointBySN(dev_sn) if not point_table then log.warn("port2pp", "not found point table with sn:",dev_sn) return end for i,v in pairs(point_table) do if v.read_write == "W" then for ii,vv in pairs(v.regs) do if vv.reginfo[2] == regaddr then return v.identifier,vv.reginfo[1] end end end end end local function port2pp(input) local result,dev_lev,addr,func,regaddr,regsize,data = ex_unpack(input) if result then local dev_sn = FetchSnByAddr(addr) if not dev_sn then -- log.warn("port2pp", "SN not found",addr) return end if func == 0x5 then --差错寄存器 log.warn("port2pp", "received a fail response !") return elseif func == 0x2 then --写寄存器响应 log.info("port2pp","write register response.",string.format("regaddr:0x%04X,regsize:%d",regaddr,regsize)) local identifier,key = findIdentifierByregaddr(dev_sn,regaddr) --TODO 需要合并重复实现 if identifier then local out_frames local regs = {} table.insert(regs,{key=key}) if not out_frames then out_frames = {} end table.insert(out_frames,{regs=regs}) internal_api.write_response(gw_sn,dev_sn,gw_port,identifier, {frames = out_frames}) end elseif func == 0x4 then --读寄存器响应 log.debug("port2pp","read register response.",string.format("regaddr:0x%04X,regsize:%d",regaddr,regsize)) else log.info("port2pp","function code invalid",func) end end end local function request_msg_pack(dev_lev,addr,regaddr,regsize) local payload = ex_pack(dev_lev,addr,0x3,regaddr,regsize) if payload then -- log.debug("request msg:",dev_lev,addr,regaddr,regsize,string.toHex(payload)) return payload end end local function write_msg_pack(dev_lev,addr,regaddr,regsize,data) local payload = ex_pack(dev_lev,addr,0x1,regaddr,regsize,data) if payload then -- log.debug("request msg:",dev_lev,addr,regaddr,regsize,string.toHex(payload)) return payload end end local function period_to_485cache(dev_lev,gw_sn,port,addr,regaddr,regsize,data) local topic = string.format( "to_dev/%s/%s/data/binary_data",gw_sn,port) local payload = ex_pack(dev_lev,addr,0x6,regaddr,regsize,data) log.debug("set period cache:",topic,string.toHex(payload)) client:publish(topic, payload) end local function value_from_msg(msg,frames,key) if not msg then log.warn("value_from_msg","Invalid paramters,",msg,frames,key) return end for _,msg_frame in ipairs(msg) do for _,reg in ipairs(msg_frame.regs) do if reg.key == key then return reg.value end end end return end local function reg_from_cache(term_addr,reg_addr,len,format,scale) local out = shmem.shmget(term_addr,reg_addr,len) if not out then log.warn(nil, "return nil as get data form memory",string.format("term_addr %d startaddr 0x%04X format:%s len:%d",term_addr,reg_addr,format,len)) return end local _,value= bunpack(out,format) if not value then log.warn(nil, "value cannot convert in batch process.",string.format("startaddr 0x%04X format:%s len:%d",reg_addr,format,len)) return end if scale then return value * scale else return value end end local function to_dev_mq_insert(msg) table.insert(to_dev_mq,msg) end -- 周期检查 local function sync_cache_poll() if #to_dev_mq >= 1 then while true do local msg = table.remove(to_dev_mq,1) if msg then client:publish(to_dev_topic, msg) else break end end end end local function parse_single_point(dev_sn,identifier,mi,frames,addr,dev_lev,v,dev_func_expr,dev_enumstr_expr) local frame_id = nil local regs = {} local obj = {} local levels = {} local extrainfos = {} if v.identifier and v.identifier ~= identifier then return end if not v.regs then log.warn(nil,"regs not found in service") return end local temp = {} if v.read_write == "R" then if v.frame_id then frame_id = v.frame_id elseif type(v.frames) == "table" then frame_id = v.frames[1] end local request_regs={} for _,vv in ipairs(v.regs) do local key = vv.reginfo[1] local reg_addr = vv.reginfo[2] local format = vv.reginfo[3] local scale = vv.reginfo[4] local enumstr = vv.reginfo[5] local batch = vv.batch local level = vv.level local extrainfo = vv.extrainfo if not format then format = ">s" end -- 寄存器地址必须用 format信息 local len = FormatLenght(format) if reg_addr and type(reg_addr) == "number" then -- 只关心直接地址,字符串格式的引用地址忽略 if not vv.batch then -- 单个寄存器 local value = reg_from_cache(addr,reg_addr,len,format,scale) if value then if not enumstr then -- 无需要转换枚举类型 if type(key) == "string" then -- 只有key是字符串的才是属性,其他是报警bit,或者nil(做空轮询用) table.insert(regs,{key=key,value=value}) elseif type(key) == "table" then for pos,sub_key in pairs(key) do if not pos or pos <= 0 then log.warn(nil, "bunpack size error",format) return end local sub_val = bit.band(bit.rshift(value,pos - 1),1) if sub_key then if type(value) == "number" then table.insert(regs,{key=sub_key,value=sub_val}) else table.insert(regs,{key=sub_key,value="--"}) end end if sub_key then obj[sub_key] = sub_val end if sub_key and level then levels[sub_key] = level end end if extrainfo then -- 存在额外信息 for _,extra in ipairs(extrainfo) do if extra[1] and extra[2] then -- 转换成K:V方便快速索引 extrainfos[extra[1] ] = extra[2] end end end end else -- 需要转换枚举类型 local result,data = enum2expr(dev_enumstr_expr,enumstr,key,value) if result and data and type(data) == "string" then log.debug(nil,"key,enumstr,data=",key,enumstr,data) table.insert(regs,{key=key,value=data}) else log.warn(nil, string.format("Value:%s not match in table list %s",value,enumstr) ) end end end else -- 连续寄存器批量操作 if type(key) == "string" then -- 只有key是字符串的才是属性,其他是报警bit,或者nil(做空轮询用) local values = {} -- 批量下发需要循环 -- log.debug(nil,string.format("======== reg: %04X", reg_addr),batch[1],batch[2]) local size = batch[1] or 1 local step = batch[2] or 1 for i = 1,size,step do local currreg = reg_addr + i-1 local value = reg_from_cache(addr,currreg,len,format,scale) if value then table.insert(values , value) end end table.insert(regs,{key=key,value=values}) -- 把数组对象插入reg end end elseif reg_addr and type(reg_addr) == "string" then -- 使用表达式引用的地址 -- 优化loadstring 只需要loadstring 一次,每个设备格式话的字符串是不一样的,所以按照设备存储 if not dev_func_expr[reg_addr] then ---- 这里reg_addr 是一个表达式,首次操作需要做一个loadstring,以后直接调用函数,目的加速调用;dev_func_expr[]是function expr2function(addr,dev_func_expr,reg_addr,key) end local f = dev_func_expr[reg_addr] if f then assert(f) local result,data = pcall(f) if result and data then log.debug(nil,"key,reg_addr,data=",key,reg_addr,data) table.insert(regs,{key=key,value=data}) else log.warn(nil, string.format("Do %s expr. failed for key:%s",key,key) ,result,data) end end end end -- 根据读写分类调用底层接口 if request_regs and #request_regs > 0 then temp = {} table.sort(request_regs) for _,reg_addr in ipairs(request_regs) do -- 新的地址开头 if not temp.regaddr then temp.regaddr = reg_addr temp.regsize = 0 temp.reg_data = '' end -- 判断地址是否连续 if (reg_addr - temp.regaddr ~= temp.regsize) then log.debug(nil,string.format("singal request addrstart:%04X length:%d",temp.regaddr,temp.regsize)) to_dev_mq_insert(request_msg_pack(dev_lev,addr,temp.regaddr,temp.regsize)) --缓存状态重新赋值 temp.regaddr = reg_addr temp.regsize = 0 temp.reg_data = '' end temp.regsize = temp.regsize + 1 end if temp.regsize > 0 then log.debug(nil,string.format("singal request addrstart:%04X length:%d",temp.regaddr,temp.regsize)) to_dev_mq_insert(request_msg_pack(dev_lev,addr,temp.regaddr,temp.regsize)) end end elseif v.read_write == "W" then for _,msg_frame in ipairs(frames) do for _,reg in ipairs(msg_frame.regs) do for ii,vv in ipairs(v.regs) do -- 新的地址开头 local key = vv.reginfo[1] local scale = vv.reginfo[4] if key == reg.key then temp.regaddr = vv.reginfo[2] temp.regsize = 0 temp.reg_data = '' if not vv.batch then -- 算出寄存器大小 local utilsize = FormatLenght(vv.reginfo[3])/2 temp.regsize = temp.regsize + utilsize -- 算出数据 local value = value_from_msg(frames,v.frames,key) if not value then log.warn(nil,"can not find the key from message,key:",key) return end if scale then if type(value) ~= "number" then log.warn(nil, "value was except as a number type.") return end temp.reg_data = temp.reg_data .. bpack(vv.reginfo[3],value/scale) else temp.reg_data = temp.reg_data .. bpack(vv.reginfo[3],value) end to_dev_mq_insert(write_msg_pack(dev_lev,addr,temp.regaddr,temp.regsize,temp.reg_data)) else local size = vv.batch[1] local step = vv.batch[2] local value = value_from_msg(frames,v.frames,key) -- log.warn(nil,"batch ----->write") if not value and type(value) ~= "table" then log.warn(nil,"[batch]can not find the key from message,key:",key) return end if size == #value then if step == 1 then temp.regaddr = vv.reginfo[2] temp.regsize = size temp.reg_data = '' temp.reg_data = temp.reg_data .. bpack(vv.reginfo[3],unpack(value,1)) log.info(nil,string.toHex(temp.reg_data)) to_dev_mq_insert(write_msg_pack(dev_lev,addr,temp.regaddr,temp.regsize,temp.reg_data)) else local format_t = string.split(reg.reginfo[3],"/") if size == #format_t then for i=0,size*(step)-step,step do temp.regaddr = vv.reginfo[2]+i local utilsize = FormatLenght(format_t[i/step+1])/2 temp.regsize = utilsize temp.reg_data = '' temp.reg_data = temp.reg_data .. bpack(format_t[i/step+1],value[i/step+1]) to_dev_mq_insert(write_msg_pack(dev_lev,addr,temp.regaddr,temp.regsize,temp.reg_data)) end else log.warn(nil,"[batch]format error,key:",key,"size:",size,"format_size:",#format_t) return end end else log.warn(nil,"[batch]format error,key:",key,"size:",size,"tavlesize:",#value) return end end break end end end end end return frame_id,regs,obj,levels,extrainfos end local function real_data_get(dev_sn,identifier,mi,frames) local addrs,dev_lev = FetchAddrBySN(dev_sn) if not addrs then log.warn(nil,"Can not found term_addr",dev_sn) return nil end if not dev_lev then log.warn(nil,"Can not found dev_lev",dev_lev) return nil end local addr = tonumber(addrs,10) -- 找服务1合并报文2找cache里的值返回 -- v 每个配置的iterator local point_table = FetchDevPointBySN(dev_sn) if point_table then local out_frames local privite_frames if not func_expr[dev_sn] then func_expr[dev_sn] = {} end local dev_func_expr = func_expr[dev_sn] if not enumstr_expr[dev_sn] then enumstr_expr[dev_sn] = {} end local dev_enumstr_expr = enumstr_expr[dev_sn] for _,v in pairs(point_table) do local frame_id,regs,objs,levels,extrainfos = parse_single_point(dev_sn,identifier,mi,frames,addr,dev_lev,v,dev_func_expr,dev_enumstr_expr) if frame_id and regs and next(regs) ~= nil then if not out_frames then out_frames = {} end if not privite_frames then privite_frames = {} end table.insert(out_frames,{frame_id=frame_id,regs=regs}) -- 对外结构,格式structrue存储 table.insert(privite_frames,{objs=objs,levels=levels,extrainfos=extrainfos}) -- 内部结构,格式K:V存储 end end if privite_frames then if (identifier == "RackAlarmRequest" or identifier == "SystemAlarmRequest") and date_synced then local immediate_changed = nil local history_changed = nil local rack_id = tonumber(string.sub(dev_sn,-2,-1)) - 1 for _,frame in ipairs(privite_frames) do local objs = frame.objs local levels = frame.levels local extrainfos = frame.extrainfos for alarm_key,value in pairs (objs) do -- 从后向前查,后面权重高 if not pack_alarm_obj[rack_id] then pack_alarm_obj[rack_id]= {} end local index,found for i,v in ipairs(pack_alarm_obj[rack_id]) do if v == alarm_key then index = i found = true break end end -- 以前没报警 and 报警发生 if not found and value == 1 then local extra = nil if extrainfos then -- 解析扩展信息 local extrainfo_key = extrainfos[alarm_key] if extrainfo_key then if not dev_func_expr[extrainfo_key] then ---- 这里extrainfo_key 是一个表达式,首次操作需要做一个loadstring,以后直接调用函数,目的加速调用;dev_func_expr[]是function expr2function(addr,dev_func_expr,extrainfo_key,alarm_key) end local f = dev_func_expr[extrainfo_key] log.info(nil,alarm_key,extrainfo_key,f) assert(f) local result,data = pcall(f) if result and data and type(data) == "string" then log.info(nil,alarm_key,data) extra = data end end end -- 保存状态 log.warn(nil,string.format("alarm[%s] flag set",levels[alarm_key],alarm_key,extra)) table.insert(pack_alarm_obj[rack_id],alarm_key) db_alarm_status_insert(stmt_insert_alarm,dev_sn,rack_id,index,alarm_key,levels[alarm_key],extra) immediate_changed = true end if found and value == 0 then -- 删除缓存 log.warn(nil,string.format("alarm[%s] flag clear",alarm_key)) table.remove(pack_alarm_obj[rack_id],index) -- 更新清除报警时间 db_alarm_status_update(stmt_update_alarm,dev_sn,rack_id,index,alarm_key) immediate_changed = true history_changed = true end end end if immediate_changed then immediate_changed = false replyImmediateAlarm(gw_sn,cluster_sn,gw_port,"ImmediateAlarm",nil,0) end if history_changed then history_changed = false replyHistoryAlarm(gw_sn,cluster_sn,gw_port,"HistoryAlarm",nil,0) end end end if out_frames then if identifier == "RealtimeGet" then local frame = out_frames[1] local obj={} for _,reg in ipairs(frame.regs) do obj[reg.key] = reg["value"] end sync_system_time(obj) end internal_api.service_response(gw_sn,dev_sn,gw_port,identifier, {frames = out_frames} ,mi) end end end --down add by zgx for system parameter option 2021-10-21 function filecopy(src,dst) os.execute("cp "..src.." "..dst) end function fileremove(src) os.execute("rm -f "..src) end local function sysParaSave() local outstr = cjson.encode(system_parameter) file = io.open(CFG_FILE..".back", "w") io.output(file) io.write(outstr) io.close(file) filecopy(CFG_FILE..".back",CFG_FILE) fileremove(CFG_FILE..".back") end local function sysParaInit() local str = "" file = io.open(CFG_FILE, "r") if file ~= nil then io.input(file) local v = 1 while v == 1 do local res = io.read() if res == nil then v=0 break end res = string.gsub(res,"^[%s\n\r\t]*(.-)[%s\n\r\t]*$","%1") str = str..res end io.close(file) end if #str == 0 then system_parameter = { identifier="systemParameter", frames = "frames", regs= { racknum=DEF_RACK_NUM, packnum=DEF_PACK_NUM, batterynum=DEF_BATT_NUM, tempsamplingnum=DEF_TEMP_NUM, batterytempsamplingnum=DEF_BATTPOLE_NUM } } else system_parameter = cjson.decode(str) end if system_parameter.regs then local cfg = system_parameter.regs racknum = cfg.racknum end end local function systemParameterProcess(dev_sn,identifier,mi,frames,cmd) local out_frames = {} local regs = {} if cmd == OPTION_WRITE then for _, msg_frame in pairs(frames) do for _,v in ipairs(msg_frame.regs) do if v.key == "racknum" then system_parameter.regs.racknum = v.value racknum = system_parameter.regs.racknum elseif v.key == "packnum" then system_parameter.regs.packnum = v.value elseif v.key == "batterynum" then system_parameter.regs.batterynum = v.value elseif v.key == "tempsamplingnum" then system_parameter.regs.tempsamplingnum = v.value elseif v.key == "batterytempsamplingnum" then system_parameter.regs.batterytempsamplingnum = v.value end table.insert(regs,{key=v.key}) end end sysParaSave() table.insert(out_frames,{frame_id=system_parameter.frames,regs=regs}) internal_api.write_response(gw_sn,dev_sn,gw_port,identifier,{frames = out_frames},mi) else regs = { {key="racknum",value=system_parameter.regs.racknum}, {key="packnum",value=system_parameter.regs.packnum}, {key="batterynum",value=system_parameter.regs.batterynum}, {key="tempsamplingnum",value=system_parameter.regs.tempsamplingnum}, {key="batterytempsamplingnum",value=system_parameter.regs.batterytempsamplingnum}, } table.insert(out_frames,{frame_id=system_parameter.frames,regs=regs}) internal_api.service_response(gw_sn,dev_sn,gw_port,identifier,{frames = out_frames},mi) end end local function realtimeTimeSetProcess(frames) local regs = {} for _,v in ipairs(frames[1].regs) do if v.key == "year" then regs.year = v.value elseif v.key == "month" then regs.month = v.value elseif v.key == "day" then regs.day = v.value elseif v.key == "hour" then regs.hour = v.value elseif v.key == "minit" then regs.minit = v.value elseif v.key == "second" then regs.second = v.value end end sync_system_time(regs) end --up add by zgx for system parameter option 2021-10-21 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 elseif obj.identifier == "ImmediateAlarm" then replyImmediateAlarm(gw_sn,dev_sn,gw_port,obj.identifier,param,obj.mi) return elseif obj.identifier == "HistoryAlarm" then replyHistoryAlarm(gw_sn,dev_sn,gw_port,obj.identifier,param,obj.mi,obj.page,obj.size) return elseif obj.identifier == "AlarmCount" then replyAlarmCount(gw_sn,dev_sn,gw_port,obj.identifier,param,obj.mi) return elseif obj.identifier == "CleanAlarm" then replyCleanAlarm() return elseif obj.identifier == "systemParameterRequest" then systemParameterProcess(dev_sn,identifier,obj.mi,obj.frames,OPTION_READ) return elseif obj.identifier == "systemParameterConfig" then systemParameterProcess(dev_sn,identifier,obj.mi,obj.frames,OPTION_WRITE) return else if obj.identifier == "RealtimeSet" then realtimeTimeSetProcess(obj.frames) end real_data_get(dev_sn,obj.identifier,obj.mi,obj.frames) end return -1, nil, 0 end local function period_report_check(sn,identifier) for k,v in ipairs(period_list)do if v then if ((os.difftime(os.time(),v.last_report) >= v.intv_report) or (os.difftime(os.time(),v.last_report) < 0)) then v.last_sample = os.time() return true end end end end local function cache_period_service() for _,node in ipairs(nodes) do local addrs,dev_lev = FetchAddrBySN(node.sn) if not addrs then log.warn("cache_period_service","Can not found term_addr") return nil end local addr = tonumber(addrs,10) if addr and addr > racknum then log.info(nil,string.format("addr %d and behind it be discard",addr)) break end local scan_regs = FetchScanRegBySN(node.sn) for _,block in ipairs(scan_regs) do if block.desc and block.start and block.stop and block.period then log.info("cache_period_service",string.format("dev:%s period describe is [%s] addr 0x%04X-0x%04X with (%d)mS",node.sn,block.desc,block.start,block.stop,block.period)) period_to_485cache(dev_lev,gw_sn,gw_port,addr,block.start,block.stop + 1 - block.start,bpack(">S",block.period)) end end end end local 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(iot2pp,payload,#payload) if not ret then log.warn("mqtt.recv",err) end elseif string.find(topic,"data/binary_data") then --二进制数据流解析 local ret,err = pcall(port2pp,payload) if not ret then log.warn("mqtt.recv",err) end -- elseif string.find(topic,"ipc/system/actions/") then -- if string.find(topic,"update_config") then -- os.execute("/app/hmi4rack/update_config.sh") -- end elseif string.find(topic,"xieneng_can/start") then log.warn("mqtt.recv","xieneng_can be restarted!!!") cache_period_service() sys.timerStart(real_data_get ,1000,cluster_sn,"RealtimeGet") else log.info("due message:",topic, payload ) end end local function on_connect() cache_period_service() -- xieneng 周期 sys.timerLoopStart(sync_cache_poll, 200) -- 消费者 real_data_get(cluster_sn,"RealtimeGet") end local function handler() local pid, status, code = syswait.wait(-1, syswait.WNOHANG) log.error('system exit', pid, status, code) --打印不出来 os.exit(0) end local function main(nodes_cfg,sn) log.info("main_loop", "=============== XIENENG BACKSTAGE RUN 1.0 ============") sysParaInit() gw_sn = sn if not gw_sn then gw_sn = io.getGatewayID() -- gw_sn = "21012C00FFFF" 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 "")) end for _,node in ipairs (nodes) do if node and node.term_addr then if node.term_addr > racknum then break end shmem.connect(node.term_addr) node.sn = string.format("%s%02d",gw_sn,node.term_addr+1) end end siganl.signal(siganl.SIGKILL,handler) siganl.signal(siganl.SIGHUP,handler) to_dev_topic = string.format( "to_dev/%s/%s/data/binary_data",gw_sn,gw_port) local subscribe={string.format("dev/+/%s/data/binary_data",gw_port),"ipc/+/+/device/+/data/Set_Rglt/+","ipc/+/process/xieneng_can/start","ipc/+/process/xieneng_can/start","ipc/system/actions/#"} local client = internal_api.mqtt_session("localhost",1883,nil,nil,subscribe,on_mqtt_msg,on_connect) identifier_index_init() log.info("main_loop", "Identifier index table initalized") db_init(DB_FILE) cluster_sn = string.format("%s%02d",gw_sn,1) sys.taskInit(function() local last_time = os.time() while true do -- 等待RTC就绪 sys.wait(1000) if (os.difftime(os.time(),last_time) >= 3600*24) or (os.difftime(os.time(),last_time) < 0) or (not date_synced and (os.difftime(os.time(),last_time) >= 5)) then last_time = os.time() real_data_get(cluster_sn,"RealtimeGet") end end end) sys.taskInit(function() local last_time = os.time() while true do for _,node in ipairs (nodes) do sys.wait(100) if node and node.term_addr > racknum then break end -- 超过配置个数不读取 if node and node.dev_lev and node.term_addr then if node.dev_lev == 2 then real_data_get(node.sn,"RackAlarmRequest") elseif node.dev_lev == 3 then real_data_get(node.sn,"SystemAlarmRequest") end end end end end) sys.taskInit(function() local report_time = 0 while true do sys.wait(1000) if ((os.difftime(os.time(),report_time) >= 3600*24*7) or (os.difftime(os.time(),report_time) < 0)) then report_time = os.time() db_autoclean() end end end) sys.taskInit(function() while true do sys.wait(200) os.execute("echo 1 > /sys/class/leds/led-heart/brightness") sys.wait(200) os.execute("echo 0 > /sys/class/leds/led-heart/brightness") end end) --启动系统框架 sys.init(0, 0) sys.run() end main(arg[1],arg[2])