From ca5139929b8b5853229ca3d63e2bca1ce82fa0ab Mon Sep 17 00:00:00 2001
From: songww <songww@inspur.com>
Date: 星期三, 13 五月 2020 14:55:06 +0800
Subject: [PATCH] 尝试修复catalog获取失败。服务重启后设备未注册仍上报keeplive处理

---
 src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java |  256 ++++++++++++++++++++++++++++++++++++++++-----------
 1 files changed, 200 insertions(+), 56 deletions(-)

diff --git a/src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java b/src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java
index 9a5551b..8a7c6cf 100644
--- a/src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java
+++ b/src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java
@@ -2,8 +2,10 @@
 
 import java.io.ByteArrayInputStream;
 import java.text.ParseException;
+import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.Iterator;
+import java.util.List;
 import java.util.Map;
 
 import javax.sip.InvalidArgumentException;
@@ -17,6 +19,8 @@
 import org.dom4j.DocumentException;
 import org.dom4j.Element;
 import org.dom4j.io.SAXReader;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Component;
 
@@ -24,11 +28,18 @@
 import com.genersoft.iot.vmp.gb28181.SipLayer;
 import com.genersoft.iot.vmp.gb28181.bean.Device;
 import com.genersoft.iot.vmp.gb28181.bean.DeviceChannel;
+import com.genersoft.iot.vmp.gb28181.bean.RecordInfo;
+import com.genersoft.iot.vmp.gb28181.bean.RecordItem;
+import com.genersoft.iot.vmp.gb28181.event.DeviceOffLineDetector;
 import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
+import com.genersoft.iot.vmp.gb28181.transmit.callback.DeferredResultHolder;
+import com.genersoft.iot.vmp.gb28181.transmit.callback.RequestMessage;
 import com.genersoft.iot.vmp.gb28181.transmit.cmd.impl.SIPCommander;
 import com.genersoft.iot.vmp.gb28181.transmit.request.ISIPRequestProcessor;
+import com.genersoft.iot.vmp.gb28181.utils.DateUtil;
 import com.genersoft.iot.vmp.gb28181.utils.XmlUtil;
 import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
+import com.genersoft.iot.vmp.utils.redis.RedisUtil;
 
 /**    
  * @Description:MESSAGE璇锋眰澶勭悊鍣�
@@ -37,7 +48,9 @@
  */
 @Component
 public class MessageRequestProcessor implements ISIPRequestProcessor {
-
+	
+	private final static Logger logger = LoggerFactory.getLogger(MessageRequestProcessor.class);
+	
 	private ServerTransaction transaction;
 	
 	private SipLayer layer;
@@ -50,6 +63,27 @@
 	
 	@Autowired
 	private EventPublisher publisher;
+	
+	@Autowired
+	private RedisUtil redis;
+	
+	@Autowired
+	private DeferredResultHolder deferredResultHolder;
+	
+	@Autowired
+	private DeviceOffLineDetector offLineDetector;
+	
+	private final static String CACHE_RECORDINFO_KEY = "CACHE_RECORDINFO_";
+	
+	private static final String MESSAGE_CATALOG = "Catalog";
+	private static final String MESSAGE_DEVICE_INFO = "DeviceInfo";
+	private static final String MESSAGE_KEEP_ALIVE = "Keepalive";
+	private static final String MESSAGE_ALARM = "Alarm";
+	private static final String MESSAGE_RECORD_INFO = "RecordInfo";
+//	private static final String MESSAGE_BROADCAST = "Broadcast";
+//	private static final String MESSAGE_DEVICE_STATUS = "DeviceStatus";
+//	private static final String MESSAGE_MOBILE_POSITION = "MobilePosition";
+//	private static final String MESSAGE_MOBILE_POSITION_INTERVAL = "Interval";
 	
 	/**   
 	 * 澶勭悊MESSAGE璇锋眰
@@ -65,17 +99,63 @@
 		this.transaction = transaction;
 		
 		Request request = evt.getRequest();
-		
-		if (new String(request.getRawContent()).contains("<CmdType>Keepalive</CmdType>")) {
-			processMessageKeepAlive(evt);
-		} else if (new String(request.getRawContent()).contains("<CmdType>Catalog</CmdType>")) {
-			processMessageCatalogList(evt);
-		} else if (new String(request.getRawContent()).contains("<CmdType>DeviceInfo</CmdType>")) {
-			processMessageDeviceInfo(evt);
-		} else if (new String(request.getRawContent()).contains("<CmdType>Alarm</CmdType>")) {
-			processMessageAlarm(evt);
+		SAXReader reader = new SAXReader();
+		Document xml;
+		try {
+			xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
+			Element rootElement = xml.getRootElement();
+			String cmd = rootElement.element("CmdType").getStringValue();
+			
+			if (MESSAGE_KEEP_ALIVE.equals(cmd)) {
+				logger.info("鎺ユ敹鍒癒eepAlive娑堟伅");
+				processMessageKeepAlive(evt);
+			} else if (MESSAGE_CATALOG.equals(cmd)) {
+				logger.info("鎺ユ敹鍒癈atalog娑堟伅");
+				processMessageCatalogList(evt);
+			} else if (MESSAGE_DEVICE_INFO.equals(cmd)) {
+				logger.info("鎺ユ敹鍒癉eviceInfo娑堟伅");
+				processMessageDeviceInfo(evt);
+			} else if (MESSAGE_ALARM.equals(cmd)) {
+				logger.info("鎺ユ敹鍒癆larm娑堟伅");
+				processMessageAlarm(evt);
+			} else if (MESSAGE_RECORD_INFO.equals(cmd)) {
+				logger.info("鎺ユ敹鍒癛ecordInfo娑堟伅");
+				processMessageRecordInfo(evt);
+			}
+		} catch (DocumentException e) {
+			e.printStackTrace();
 		}
 		
+	}
+	
+	/**
+	 * 鏀跺埌deviceInfo璁惧淇℃伅璇锋眰 澶勭悊
+	 * @param evt
+	 */
+	private void processMessageDeviceInfo(RequestEvent evt) {
+		try {
+			Element rootElement = getRootElement(evt);
+			Element deviceIdElement = rootElement.element("DeviceID");
+			String deviceId = deviceIdElement.getText().toString();
+			
+			Device device = storager.queryVideoDevice(deviceId);
+			if (device == null) {
+				return;
+			}
+			device.setName(XmlUtil.getText(rootElement,"DeviceName"));
+			device.setManufacturer(XmlUtil.getText(rootElement,"Manufacturer"));
+			device.setModel(XmlUtil.getText(rootElement,"Model"));
+			device.setFirmware(XmlUtil.getText(rootElement,"Firmware"));
+			storager.update(device);
+			
+			RequestMessage msg = new RequestMessage();
+			msg.setDeviceId(deviceId);
+			msg.setType(DeferredResultHolder.CALLBACK_CMD_DEVICEINFO);
+			msg.setData(device);
+			deferredResultHolder.invokeResult(msg);
+		} catch (DocumentException e) {
+			e.printStackTrace();
+		}
 	}
 	
 	/***
@@ -84,11 +164,7 @@
 	 */
 	private void processMessageCatalogList(RequestEvent evt) {
 		try {
-			Request request = evt.getRequest();
-			SAXReader reader = new SAXReader();
-			reader.setEncoding("GB2312");
-			Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
-			Element rootElement = xml.getRootElement();
+			Element rootElement = getRootElement(evt);
 			Element deviceIdElement = rootElement.element("DeviceID");
 			String deviceId = deviceIdElement.getText().toString();
 			Element deviceListElement = rootElement.element("DeviceList");
@@ -152,36 +228,12 @@
 				}
 				// 鏇存柊
 				storager.update(device);
+				RequestMessage msg = new RequestMessage();
+				msg.setDeviceId(deviceId);
+				msg.setType(DeferredResultHolder.CALLBACK_CMD_CATALOG);
+				msg.setData(device);
+				deferredResultHolder.invokeResult(msg);
 			}
-		} catch (DocumentException e) {
-			e.printStackTrace();
-		}
-	}
-	
-	/***
-	 * 鏀跺埌deviceInfo璁惧淇℃伅璇锋眰 澶勭悊
-	 * @param evt
-	 */
-	private void processMessageDeviceInfo(RequestEvent evt) {
-		try {
-			Request request = evt.getRequest();
-			SAXReader reader = new SAXReader();
-			// reader.setEncoding("GB2312");
-			Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
-			Element rootElement = xml.getRootElement();
-			Element deviceIdElement = rootElement.element("DeviceID");
-			String deviceId = deviceIdElement.getText().toString();
-			
-			Device device = storager.queryVideoDevice(deviceId);
-			if (device == null) {
-				return;
-			}
-			device.setName(XmlUtil.getText(rootElement,"DeviceName"));
-			device.setManufacturer(XmlUtil.getText(rootElement,"Manufacturer"));
-			device.setModel(XmlUtil.getText(rootElement,"Model"));
-			device.setFirmware(XmlUtil.getText(rootElement,"Firmware"));
-			storager.update(device);
-			cmder.catalogQuery(device);
 		} catch (DocumentException e) {
 			e.printStackTrace();
 		}
@@ -193,11 +245,7 @@
 	 */
 	private void processMessageAlarm(RequestEvent evt) {
 		try {
-			Request request = evt.getRequest();
-			SAXReader reader = new SAXReader();
-			// reader.setEncoding("GB2312");
-			Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
-			Element rootElement = xml.getRootElement();
+			Element rootElement = getRootElement(evt);
 			Element deviceIdElement = rootElement.element("DeviceID");
 			String deviceId = deviceIdElement.getText().toString();
 			
@@ -222,18 +270,114 @@
 	 */
 	private void processMessageKeepAlive(RequestEvent evt){
 		try {
+			Element rootElement = getRootElement(evt);
+			String deviceId = XmlUtil.getText(rootElement,"DeviceID");
 			Request request = evt.getRequest();
-			Response response = layer.getMessageFactory().createResponse(Response.OK,request);
-			SAXReader reader = new SAXReader();
-			Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
-			// reader.setEncoding("GB2312");
-			Element rootElement = xml.getRootElement();
-			Element deviceIdElement = rootElement.element("DeviceID");
+			Response response = null;
+			if (offLineDetector.isOnline(deviceId)) {
+				response = layer.getMessageFactory().createResponse(Response.OK,request);
+				publisher.onlineEventPublish(deviceId, VideoManagerConstants.EVENT_ONLINE_KEEPLIVE);
+			} else {
+				response = layer.getMessageFactory().createResponse(Response.BAD_REQUEST,request);
+			}
 			transaction.sendResponse(response);
-			publisher.onlineEventPublish(deviceIdElement.getText(), VideoManagerConstants.EVENT_ONLINE_KEEPLIVE);
 		} catch (ParseException | SipException | InvalidArgumentException | DocumentException e) {
 			e.printStackTrace();
 		}
 	}
+	
+	/***
+	 * 鏀跺埌catalog璁惧鐩綍鍒楄〃璇锋眰 澶勭悊
+	 * TODO 杩囨湡鏃堕棿鏆傛椂鍐欐180绉掞紝鍚庣画涓嶥eferredResult瓒呮椂鏃堕棿淇濇寔涓�鑷�
+	 * @param evt
+	 */
+	private void processMessageRecordInfo(RequestEvent evt) {
+		try {
+			RecordInfo recordInfo = new RecordInfo();
+			Element rootElement = getRootElement(evt);
+			Element deviceIdElement = rootElement.element("DeviceID");
+			String deviceId = deviceIdElement.getText().toString();
+			recordInfo.setDeviceId(deviceId);
+			recordInfo.setName(XmlUtil.getText(rootElement,"Name"));
+			recordInfo.setSumNum(Integer.parseInt(XmlUtil.getText(rootElement,"SumNum")));
+			String sn = XmlUtil.getText(rootElement,"SN");
+			Element recordListElement = rootElement.element("RecordList");
+			if (recordListElement == null) {
+				return;
+			}
+			
+			Iterator<Element> recordListIterator = recordListElement.elementIterator();
+			List<RecordItem> recordList = new ArrayList<RecordItem>();
+			if (recordListIterator != null) {
+				RecordItem record = new RecordItem();
+				// 閬嶅巻DeviceList
+				while (recordListIterator.hasNext()) {
+					Element itemRecord = recordListIterator.next();
+					Element recordElement = itemRecord.element("DeviceID");
+					if (recordElement == null) {
+						continue;
+					}
+					record = new RecordItem();
+					record.setDeviceId(XmlUtil.getText(itemRecord,"DeviceID"));
+					record.setName(XmlUtil.getText(itemRecord,"Name"));
+					record.setFilePath(XmlUtil.getText(itemRecord,"FilePath"));
+					record.setAddress(XmlUtil.getText(itemRecord,"Address"));
+					record.setStartTime(DateUtil.ISO8601Toyyyy_MM_dd_HH_mm_ss(XmlUtil.getText(itemRecord,"StartTime")));
+					record.setEndTime(DateUtil.ISO8601Toyyyy_MM_dd_HH_mm_ss(XmlUtil.getText(itemRecord,"EndTime")));
+					record.setSecrecy(itemRecord.element("Secrecy") == null? 0:Integer.parseInt(XmlUtil.getText(itemRecord,"Secrecy")));
+					record.setType(XmlUtil.getText(itemRecord,"Type"));
+					record.setRecordId(XmlUtil.getText(itemRecord,"RecorderID"));
+					recordList.add(record);
+				}
+				recordInfo.setRecordList(recordList);
+			}
+			
+			// 瀛樺湪褰曞儚涓斿鏋滃綋鍓嶅綍鍍忔槑缁嗕釜鏁板皬浜庢�绘潯鏁帮紝璇存槑鎷嗗寘杩斿洖锛岄渶瑕佺粍瑁咃紝鏆備笉杩斿洖
+			if (recordInfo.getSumNum() > 0 && recordList.size() > 0 && recordList.size() < recordInfo.getSumNum()) {
+				// 涓洪槻姝㈣繛缁姹傝璁惧鐨勫綍鍍忔暟鎹紝杩斿洖鏁版嵁閿欎贡锛岀壒澧炲姞sn杩涜鍖哄垎
+				String cacheKey = CACHE_RECORDINFO_KEY+deviceId+sn;
+				// TODO 鏆傛椂鐩存帴鎿嶄綔redis瀛樺偍锛屽悗缁皝瑁呬笓鐢ㄧ紦瀛樻帴鍙o紝鏀逛负鏈湴鍐呭瓨缂撳瓨
+				if (redis.hasKey(cacheKey)) {
+					List<RecordItem> previousList = (List<RecordItem>) redis.get(cacheKey);
+					if (previousList != null && previousList.size() > 0) {
+						recordList.addAll(previousList);
+					}
+					// 鏈垎鏀〃绀哄綍鍍忓垪琛ㄨ鎷嗗寘锛屼笖鍔犱笂涔嬪墠鐨勬暟鎹繕鏄笉澶�,淇濆瓨缂撳瓨杩斿洖锛岀瓑寰呬笅涓寘鍐嶅鐞�
+					if (recordList.size() < recordInfo.getSumNum()) {
+						redis.set(cacheKey, recordList, 180);
+						return;
+					} else {
+						// 鏈垎鏀〃绀哄綍鍍忚鎷嗗寘锛屼絾鍔犱笂涔嬪墠鐨勬暟鎹瓒冲锛岃繑鍥炲搷搴�
+						// 鍥犺澶囧績璺虫湁鐩戝惉redis杩囨湡鏈哄埗锛屼负鎻愰珮鎬ц兘锛屾澶勬墜鍔ㄥ垹闄�
+						redis.del(cacheKey);
+					}
+				} else {
+					// 鏈垎鏀湁涓ょ鍙兘锛�1銆佸綍鍍忓垪琛ㄨ鎷嗗寘锛屼笖鏄涓�涓寘,鐩存帴淇濆瓨缂撳瓨杩斿洖锛岀瓑寰呬笅涓寘鍐嶅鐞�
+					//             2銆佷箣鍓嶆湁鍖咃紝浣嗚秴鏃舵竻绌轰簡锛岄偅涔堣繖娆n鎵规鐨勫搷搴旀暟鎹凡缁忎笉瀹屾暣锛岀瓑寰呰繃鏈熸椂闂村悗redis鑷姩娓呯┖鏁版嵁
+					redis.set(cacheKey, recordList, 180);
+					return;
+				}
+				
+			}
+			// 璧板埌杩欓噷锛屾湁浠ヤ笅鍙兘锛�1銆佹病鏈夊綍鍍忎俊鎭�,绗竴娆℃敹鍒皉ecordinfo鐨勬秷鎭嵆杩斿洖鍝嶅簲鏁版嵁锛屾棤redis鎿嶄綔
+			//               2銆佹湁褰曞儚鏁版嵁锛屼笖绗竴娆″嵆鏀跺埌瀹屾暣鏁版嵁锛岃繑鍥炲搷搴旀暟鎹紝鏃爎edis鎿嶄綔
+			//	             3銆佹湁褰曞儚鏁版嵁锛屽湪瓒呮椂鏃堕棿鍐呮敹鍒板娆″寘缁勮鍚庢暟閲忚冻澶燂紝杩斿洖鏁版嵁
+			RequestMessage msg = new RequestMessage();
+			msg.setDeviceId(deviceId);
+			msg.setType(DeferredResultHolder.CALLBACK_CMD_RECORDINFO);
+			msg.setData(recordInfo);
+			deferredResultHolder.invokeResult(msg);
+		} catch (DocumentException e) {
+			e.printStackTrace();
+		}
+	}
+	
+	private Element getRootElement(RequestEvent evt) throws DocumentException {
+		Request request = evt.getRequest();
+		SAXReader reader = new SAXReader();
+		reader.setEncoding("GB2312");
+		Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
+		return xml.getRootElement();
+	}
 
 }

--
Gitblit v1.8.0