cjson=require("cjson") MQTT=require("mqtt_library") base64=require("base64") dyutils=require("dyutils") bit=require("bit") lpack=require("lua_pack") syslog=require("syslog") bcd=require("bcd") sqlite3=require("lsqlite3") bnot = bit.bnot band, bor, bxor = bit.band, bit.bor, bit.bxor lshift, rshift, rol = bit.lshift, bit.rshift, bit.rol --MQTT.Utility.set_debug(true) syslog.openlog("dy_lora upgrade", syslog.LOG_PERROR + syslog.LOG_ODELAY, "LOG_USER") broker_addr = "localhost" broker_port = 1883 once_block_max = 128 upgrade_num = 1 --升级编号 nodes_cfg = "/app/node/nodes_cfg.json" broadcat_sn = '*' process_total = 1 process_curr = 1 serial_num = 1001 time_last = 0 function on_mqtt_msg(topic, payload) -- application specific code -- print("receive ", topic, payload); payload_obj = cjson.decode(payload) if(string.find(topic,'regulate_response_lora') ~= nil)then coroutine.yield(topic, payload_obj) elseif(string.find(payload_obj['src_identifier'],'__transparent') ~= nil)then coroutine.yield(topic, payload_obj) end end function mqtt_listen() while(true) do if(time_last ~= 0 and os.time() >= (time_last + 10))then coroutine.yield(nil, nil) end mqtt_client:handler() end end function send_data_to_sn(mqtt_client, gw_sn, dev_sn, mi, data, obj_send) topic = obj_send['topic'] obj_send['mi'] = mi obj_send['len'] = #data obj_send['data_b64'] = base64.encode(data) obj_send['sn'] = dev_sn json_str = cjson.encode(obj_send) mqtt_client:publish(topic, json_str) -- print(json_str) time_last = os.time() end function send_upgrade_process(mqtt_client, gw_sn, dev_sn, mi, Upgrade_Status, Upgrade_Progress) topic = string.format("ipc/%s/upgrade/Rsp_EventData", gw_sn) local payload={} payload['identifier'] = '__Upgrade_Progress' payload['sn'] = dev_sn payload['mi'] = mi payload['time'] = os.time() local progress = {} progress['Upgrade_Status'] = Upgrade_Status progress['Upgrade_Progress'] = Upgrade_Progress payload['tags'] = progress json_str = cjson.encode(payload) -- print(json_str) mqtt_client:publish(topic, json_str) end -- send response to dev_sn -- function send_upgrade_response(mqtt_client, gw_sn, mi, Upgrade_Response) -- topic = string.format("ipc/%s/upgrade/Rsp_EventData", gw_sn) -- local payload={} -- payload['identifier'] = '__Upgrade_Response' -- payload['sn'] = dev_sn -- payload['mi'] = mi -- payload['mi'] = mi -- payload['time'] = os.time() -- local response = {} -- response['Upgrade_Response'] = Upgrade_Response -- payload['tags'] = response -- json_str = cjson.encode(payload) -- print(json_str) -- mqtt_client:publish(topic, json_str) -- end function is_block_empty(buff,size) local i=1 while(i <=size) do if(string.byte(buff,i) ~= 0xff)then break end i = i + 1 end if(i < size)then return 0 end return 1; end function upgrade_node(json_str, gw_sn, firmware, send_out_json) local obj=cjson.decode(json_str) local obj_send = cjson.decode(send_out_json) local ret = 0 mqtt_client = MQTT.client.create(broker_addr, broker_port, on_mqtt_msg) mqtt_client:connect("upgrade_dy_lora") topic = string.format("ipc/%s/+/raw_data", gw_sn) mqtt_client:subscribe({topic}) topic = string.format("ipc/%s/+/regulate_response_lora", gw_sn) mqtt_client:subscribe({topic}) mqtt = coroutine.create(mqtt_listen) -- 获取firmware的version信息 _,dev_type,firmware_version = string.match(firmware,"^(.*)_(%x+)_V(%d+)") syslog.syslog("LOG_INFO", "firmware version:"..firmware_version) syslog.syslog("LOG_INFO", "device type version:"..dev_type) -- 解bin文件计算去掉FF后包大小 file = io.open(firmware, "r") if(file == nil)then syslog.syslog("LOG_ERROR", "fetch nodes config file failed!") end local bin_buf = file:read("*a") file:seek("set") local index = 0 local total = #bin_buf blks={} while (total > 0) do local blk_size if(total >= once_block_max) then blk_size = once_block_max else blk_size = total end if(is_block_empty(string.sub(bin_buf,once_block_max*index + 1 ),blk_size) == 0)then blk={} blk['data']=string.sub(bin_buf,once_block_max*index + 1,once_block_max*index + blk_size) -- print(#blk['data'],index,blk_size,once_block_max*index + 1 + blk_size) blk['offset']=index table.insert(blks,blk) if(#blks >= 2000)then syslog.syslog("LOG_ERROR", string.format("upgread block too large %d",#blks)) end else -- print(string.format("block [%d] all 0xff",index) ) end total = total-blk_size; index = index + 1; end syslog.syslog("LOG_INFO", string.format("block split done ,size:%d",#blks)) file:close() process_total = process_total + #blks -- 读取配置文件找到匹配的SN file = io.open(nodes_cfg, "r") if(file == nil)then syslog.syslog("LOG_ERROR", "open nodes config file failed!") end local bin_buf = file:read("*a") obj_cfg = cjson.decode(bin_buf) if(obj_cfg == nil) then syslog.syslog("LOG_ERROR", "decode nodes config file failed!") return -1; end file:close() dev_sn={} for key,value in ipairs(obj_cfg['nodes_cfg']) do local sn_type = string.sub(value['sn'],5,6) if(value['connect_port']=='LORA_1' and sn_type == dev_type)then info = {} -- info['sn'] = '210126F0000F' info['sn'] = value['sn'] info['status'] = 'ready' table.insert( dev_sn, info) -- break end end for key,value in ipairs(dev_sn) do syslog.syslog("LOG_INFO", string.format("expect upgrade sequence %d device sn %s", key,value['sn'])) send_upgrade_process(mqtt_client, gw_sn, value['sn'], serial_num, value['status'], 0) serial_num = serial_num + 1 end process_total = process_total + #dev_sn -- 通知空中升级 i=0; process_total = process_total + 4 while (i < 4) do i=i+1 func_code = 0x71 block_pkt = bpack(">SCSC= 10)then break end end -- 设备逐一确认包缺失情况 for key,value in ipairs(dev_sn) do local retry = 10 while(retry > 0) do repeat syslog.syslog("LOG_INFO", string.format("device %s check upgrade,retry %d",value['sn'],retry)) retry = retry - 1 func_code = 0x73 block_pkt = bpack(">SC",0x0002,func_code); send_data_to_sn(mqtt_client, gw_sn, value['sn'], 903, block_pkt, obj_send) received = false loop = true while(loop)do repeat ret, topic, payload_obj = coroutine.resume(mqtt) if(payload_obj == nil or topic == nil)then syslog.syslog("LOG_WARNING", "miss resoponse,timeout") loop = false break end if(payload_obj == nil or payload_obj.mi == nil or payload_obj.data_b64 == nil)then syslog.syslog("LOG_WARNING", string.format("miss resoponse,slave node not ack", topic)) break end received = true loop = false until true end if(received == false)then break end local rcv_bin = base64.decode(payload_obj.data_b64) _,_,rcv_func_code,control_code,verson_curr = bunpack(rcv_bin,">SCCSC",0x00ff,func_code); send_data_to_sn(mqtt_client, gw_sn, broadcat_sn, 904, block_pkt, obj_send) ret, topic, payload_obj = coroutine.resume(mqtt) retry = 0 break else -- 解包得到丢失的包序号 miss_blks={} offset = 5 while(offset < #rcv_bin) do _,result = bunpack(string.sub(rcv_bin,offset), "SC